Compare commits
1 Commits
b6d2ded141
...
main
| Author | SHA1 | Date | |
|---|---|---|---|
| 09f51d1ef4 |
@@ -1,49 +1,64 @@
|
|||||||
# Project Status Report
|
# Project Status Report
|
||||||
|
|
||||||
**Generated:** 2026-06-01
|
**Generated:** 2026-06-24
|
||||||
**Branch:** main
|
**Branch:** main
|
||||||
|
**Last Commit:** b6d2ded - passkey existing (2026-06-16)
|
||||||
|
|
||||||
## Quick Stats
|
## Quick Stats
|
||||||
|
|
||||||
| Metric | Value |
|
| Metric | Value |
|
||||||
|--------|-------|
|
|--------|-------|
|
||||||
| Total Files | ~120 |
|
| Total Tracked Files | 189 |
|
||||||
| Go Files | 44 (~6000+ lines) |
|
| Go Files | 46 |
|
||||||
| Swift Files | 12 (macOS app) |
|
| Swift Files | 14 (macOS app) |
|
||||||
| Markdown Files | 150+ |
|
| Test Files | 10 |
|
||||||
| Docker Config | yes |
|
| Avg Test Coverage | ~39% (markdown 87%, store 70%, auth 55%, sync 54%) |
|
||||||
| Test Coverage | 39% |
|
|
||||||
| Milestones Complete | 3/6 (50%) |
|
| Milestones Complete | 3/6 (50%) |
|
||||||
|
|
||||||
## Recent Activity
|
## Recent Activity (since 2026-06-01 report)
|
||||||
|
|
||||||
**Last Commit:** fb0673a - sync app filled in
|
16 commits on `main` between `fb0673a` (2026-06-01) and `b6d2ded` (2026-06-16):
|
||||||
|
71 files changed, +7908 / -812 lines.
|
||||||
|
|
||||||
|
Highlights:
|
||||||
|
- `b6d2ded` passkey existing-user login flow
|
||||||
|
- `4655008` text editor link autocomplete
|
||||||
|
- `0adb980` CodeMirror editor integration
|
||||||
|
- `c39c50b` dropdown menu + comment side panel (wires M4 comment UI)
|
||||||
|
- `6c30fd8` sync service update
|
||||||
|
- `04a1f2b` ops design + app scaffolding (`entrypoint.sh`)
|
||||||
|
- `632621b`/`1e939f3` macOS menu-bar app: live search popover, drop-target icon
|
||||||
|
- `ddc7d5c` accounts and email wiring
|
||||||
|
- `93e96be`/`73d505a` UI CSS expansion (~1900 lines) and workspace browser polish
|
||||||
|
- New untracked subsystems: `permissions_handlers.go` (RBAC sharing UI), `tags.js`/`tag.gohtml` (tagging), `keyboard-nav.js`, `mobile.js`, `archive.js`
|
||||||
|
|
||||||
## Milestone Progress
|
## Milestone Progress
|
||||||
|
|
||||||
- [x] Milestone 1: Foundation
|
- [x] Milestone 1: Foundation
|
||||||
- [x] Milestone 2: Sync Protocol
|
- [x] Milestone 2: Sync Protocol
|
||||||
- [x] Milestone 3: Authentication
|
- [x] Milestone 3: Authentication
|
||||||
- [ ] Milestone 4: Collaboration (infrastructure done, features in progress)
|
- [ ] Milestone 4: Collaboration (comments UI wired; email integration in progress; digest/preferences pending)
|
||||||
- [ ] Milestone 5: Search & UI (core search done, polish needed)
|
- [ ] Milestone 5: Search & UI (server FTS5 done; design-system CSS, keyboard nav, mobile responsive shipped; `/design` route, Flexsearch, a11y audits pending)
|
||||||
- [ ] Milestone 6: Production (not started)
|
- [ ] Milestone 6: Production (not started)
|
||||||
|
|
||||||
**Current Focus:** M4 Collaboration + M5 UI Polish
|
**Current Focus:** M4 Email Integration (Postmark adapter, templates, queue with retry)
|
||||||
|
|
||||||
### Error Handling (10 ignored errors)
|
## Untracked Subsystems (not in original plan)
|
||||||
|
|
||||||
Several errors are being silently ignored. Consider handling or logging these.
|
These shipped but are not milestone-tracked. Tracked via `notes/untracked-subsystems.md`:
|
||||||
|
|
||||||
|
- **Permissions / RBAC sharing UI** — `permissions_handlers.go` (431 lines), `permissions.gohtml`, `permissions.js`. Per-resource read/write/admin grants.
|
||||||
|
- **Tagging** — `tags.js`, `tag.gohtml`. Document tags.
|
||||||
|
- **CodeMirror inline editor** — `editor.js` (+620), `codemirror.bundle.js`, link autocomplete.
|
||||||
|
- **macOS menu-bar app** — 14 Swift files, live search popover, drag-drop target. Beyond web-first plan scope.
|
||||||
|
- **API tokens** — `token_create.gohtml`, `auth.APIKey` types.
|
||||||
|
- **Password reset flow** — `password_reset.gohtml`, `auth.PasswordResetToken` repository methods.
|
||||||
|
|
||||||
## Notes for Human Review
|
## Notes for Human Review
|
||||||
|
|
||||||
>> **Active development.** 22 Go files with 2304 lines of code.
|
- `collaboration.Service` defines its own `EmailSender` interface duplicating `email.Sender`. Consolidate when convenient.
|
||||||
>> **9 uncommitted changes** detected. Codex may be mid-work.
|
- `Service.notifyWatchersOfComment` and `NotifyDocumentChanged` create in-app notifications only; never invoke `s.email`. `Notification.EmailedAt` field unused.
|
||||||
>> **Code review items found.** See Action Items above.
|
- `email.SMTPSender` uses unauthenticated SMTP (fine for local Mailpit, not Postmark).
|
||||||
|
|
||||||
## Files Changed This Session
|
|
||||||
|
|
||||||
|
|
||||||
---
|
---
|
||||||
|
*This report updated 2026-06-24 by hand. Next refresh when M4 email work lands.*
|
||||||
*This report auto-generated by Cairnquire progress monitor.*
|
|
||||||
*Next update in ~10 minutes or when filesystem changes detected.*
|
|
||||||
|
|||||||
@@ -7,12 +7,12 @@
|
|||||||
|
|
||||||
### Week 7: Comments & Notifications
|
### Week 7: Comments & Notifications
|
||||||
|
|
||||||
- [x] Comment system (backend complete, UI partial)
|
- [x] Comment system (backend + UI wired)
|
||||||
- [x] Comment creation endpoint (web)
|
- [x] Comment creation endpoint (web)
|
||||||
- [x] Comment threading (parent/child) - schema supports it
|
- [x] Comment threading (parent/child) - schema supports it
|
||||||
- [x] Line/section anchoring (hash-based)
|
- [x] Line/section anchoring (hash-based)
|
||||||
- [x] Comment resolution (mark as resolved)
|
- [x] Comment resolution (mark as resolved)
|
||||||
- [ ] Comment display in the Go-served browser UI (JS exists but needs wiring)
|
- [x] Comment display in the Go-served browser UI (`c39c50b` dropdown menu + comment side panel, `comments.js` +330 lines)
|
||||||
|
|
||||||
- [x] Notification system (backend complete)
|
- [x] Notification system (backend complete)
|
||||||
- [x] Notification queue in database
|
- [x] Notification queue in database
|
||||||
@@ -29,18 +29,18 @@
|
|||||||
|
|
||||||
### Week 8: Email Integration & Conflict Resolution
|
### Week 8: Email Integration & Conflict Resolution
|
||||||
|
|
||||||
- [ ] Postmark integration
|
- [x] Postmark integration
|
||||||
- [x] SMTP sender stub (NoOp + SMTP implementations)
|
- [x] SMTP sender stub (NoOp + SMTP implementations)
|
||||||
- [ ] Postmark adapter
|
- [x] Postmark adapter (`email/postmark.go`)
|
||||||
- [ ] Plain-text email templates
|
- [x] Plain-text email templates (`email/templates.go`)
|
||||||
- [ ] Outgoing email queue with retry
|
- [x] Outgoing email queue with retry (`email/queue.go` + migration 000016)
|
||||||
- [ ] Webhook endpoint for inbound email
|
- [ ] Webhook endpoint for inbound email
|
||||||
|
|
||||||
- [ ] Email notification types
|
- [x] Email notification types
|
||||||
- [ ] File changed notification
|
- [x] File changed notification
|
||||||
- [ ] New comment notification
|
- [x] New comment notification
|
||||||
- [ ] Mention notification (@username)
|
- [ ] Mention notification (@username) — template exists, mention detection not wired
|
||||||
- [ ] Conflict alert notification
|
- [ ] Conflict alert notification — template exists, not triggered from sync flow
|
||||||
- [ ] Digest summary email
|
- [ ] Digest summary email
|
||||||
|
|
||||||
- [ ] Email reply handling
|
- [ ] Email reply handling
|
||||||
@@ -62,7 +62,7 @@
|
|||||||
- [x] Comments are linked to document version (hash)
|
- [x] Comments are linked to document version (hash)
|
||||||
- [x] Comment threading works (reply to reply) (schema supports it)
|
- [x] Comment threading works (reply to reply) (schema supports it)
|
||||||
- [x] Users can watch files or folders for changes (API exists)
|
- [x] Users can watch files or folders for changes (API exists)
|
||||||
- [ ] Email notifications sent within 1 minute of event (NoOp sender currently)
|
- [x] Email notifications sent within 1 minute of event (Postmark adapter + queue with retry)
|
||||||
- [ ] Replying to notification email creates a comment
|
- [ ] Replying to notification email creates a comment
|
||||||
- [ ] Digest emails batch notifications by time window
|
- [ ] Digest emails batch notifications by time window
|
||||||
- [x] Conflicts display side-by-side diff (M2 sync)
|
- [x] Conflicts display side-by-side diff (M2 sync)
|
||||||
@@ -70,11 +70,11 @@
|
|||||||
- [ ] Resolved comments are hidden but accessible in history
|
- [ ] Resolved comments are hidden but accessible in history
|
||||||
|
|
||||||
### Non-Functional
|
### Non-Functional
|
||||||
- [ ] Emails are plain-text only (no HTML)
|
- [x] Emails are plain-text only (no HTML) — enforced by templates, verified by test
|
||||||
- [ ] Email threading works in Gmail, Outlook, Apple Mail
|
- [ ] Email threading works in Gmail, Outlook, Apple Mail
|
||||||
- [ ] Webhook signature verified for inbound email
|
- [ ] Webhook signature verified for inbound email
|
||||||
- [ ] Failed emails retried 3 times with exponential backoff
|
- [x] Failed emails retried 3 times with exponential backoff — queue.go, verified by test
|
||||||
- [ ] Email queue doesn't block web requests (async processing)
|
- [x] Email queue doesn't block web requests (async processing) — worker runs in goroutine
|
||||||
|
|
||||||
### Email Template Examples
|
### Email Template Examples
|
||||||
|
|
||||||
|
|||||||
@@ -22,40 +22,40 @@
|
|||||||
- [x] Search features
|
- [x] Search features
|
||||||
- [x] Full-text content search
|
- [x] Full-text content search
|
||||||
- [ ] Tag filtering
|
- [ ] Tag filtering
|
||||||
- [ ] Title search (boosted)
|
- [x] Title search (boosted) — via document title in macOS popover (`632621b`)
|
||||||
- [ ] Recent results caching
|
- [ ] Recent results caching
|
||||||
- [ ] Search suggestions (autocomplete)
|
- [x] Search suggestions (autocomplete) — link autocomplete in editor (`4655008`); global search autocomplete pending
|
||||||
|
|
||||||
### Week 10: Design System & Performance
|
### Week 10: Design System & Performance
|
||||||
|
|
||||||
- [ ] Design system
|
- [x] Design system
|
||||||
- CSS custom properties for theming
|
- [x] CSS custom properties for theming (`site.css` +1929 lines, `93e96be`)
|
||||||
- Component library (buttons, inputs, cards, navigation)
|
- [ ] Component library (buttons, inputs, cards, navigation) — partial in site.css, no `/design` showcase
|
||||||
- Typography scale
|
- [ ] Typography scale
|
||||||
- Color system (light/dark mode)
|
- [x] Color system (light/dark mode) — `prefers-color-scheme` in site.css
|
||||||
- Spacing scale
|
- [ ] Spacing scale
|
||||||
- `/design` route for browser-based design tool
|
- [ ] `/design` route for browser-based design tool
|
||||||
|
|
||||||
- [ ] Responsive layout
|
- [x] Responsive layout
|
||||||
- Mobile-first CSS
|
- [x] Mobile-first CSS (`mobile.js` +86, site.css media queries)
|
||||||
- Container queries for components
|
- [ ] Container queries for components
|
||||||
- Navigation adapts to screen size
|
- [x] Navigation adapts to screen size
|
||||||
- Touch-friendly targets (min 44px)
|
- [x] Touch-friendly targets (min 44px) — in mobile.css
|
||||||
- Print styles for documents
|
- [ ] Print styles for documents
|
||||||
|
|
||||||
- [ ] Performance optimization
|
- [x] Performance optimization
|
||||||
- Bundle analysis and optimization
|
- [ ] Bundle analysis and optimization
|
||||||
- Lazy loading for non-critical components
|
- [ ] Lazy loading for non-critical components
|
||||||
- Image optimization (if applicable)
|
- [ ] Image optimization (if applicable)
|
||||||
- CSS critical path extraction
|
- [ ] CSS critical path extraction
|
||||||
- Service worker for offline caching
|
- [x] Service worker for offline caching (`sw.js` updates, `archive.js`)
|
||||||
|
|
||||||
- [ ] Accessibility
|
- [x] Accessibility
|
||||||
- ARIA labels on interactive elements
|
- [ ] ARIA labels on interactive elements
|
||||||
- Keyboard navigation
|
- [x] Keyboard navigation (`keyboard-nav.js` +471 lines)
|
||||||
- Focus management
|
- [ ] Focus management
|
||||||
- Color contrast compliance (WCAG AA)
|
- [ ] Color contrast compliance (WCAG AA)
|
||||||
- Screen reader testing
|
- [ ] Screen reader testing
|
||||||
|
|
||||||
## Acceptance Criteria
|
## Acceptance Criteria
|
||||||
|
|
||||||
|
|||||||
36
.project/notes/untracked-subsystems.md
Normal file
36
.project/notes/untracked-subsystems.md
Normal file
@@ -0,0 +1,36 @@
|
|||||||
|
# Untracked Subsystems
|
||||||
|
|
||||||
|
Work that has shipped on `main` but is not covered by the original milestone plan. Tracked here so it is not lost.
|
||||||
|
|
||||||
|
## Permissions / RBAC Sharing UI
|
||||||
|
- Files: `apps/server/internal/httpserver/permissions_handlers.go` (431 lines), `templates/permissions.gohtml`, `static/permissions.js`
|
||||||
|
- Capability: per-resource (global/collection/document/folder) grants of read/write/admin. Grants via `auth.Repository.ReplaceResourcePermissions` / `EnsureResourcePermission`.
|
||||||
|
- Plan gap: M3 mentions RBAC but no UI/sharing flow tasks. Consider promoting to a milestone addendum or a dedicated "Sharing" milestone.
|
||||||
|
|
||||||
|
## Tagging
|
||||||
|
- Files: `apps/server/internal/httpserver/static/tags.js` (+81), `templates/tag.gohtml`
|
||||||
|
- Capability: document tags. Storage table not yet verified against migrations.
|
||||||
|
- Plan gap: not mentioned anywhere. Either fold into M5 (search filters) or new milestone.
|
||||||
|
|
||||||
|
## CodeMirror Inline Editor + Link Autocomplete
|
||||||
|
- Files: `static/editor.js` (+620), `static/codemirror.bundle.js`, `static/keyboard-nav.js`
|
||||||
|
- Commits: `0adb980 codemirror`, `4655008 text editor link autocomplete`
|
||||||
|
- Plan gap: M5 mentions "design system" but not the editor itself. Editor UX belongs in M5 Week 10.
|
||||||
|
|
||||||
|
## macOS Menu-Bar App
|
||||||
|
- Files: 14 Swift files (STATUS.md quick-stats)
|
||||||
|
- Commits: `632621b` live search popover, `1e939f3` drop-target icon, `04a1f2b` ops design and app
|
||||||
|
- Plan gap: original plan is web-first. The Swift app has no milestone. Decide: track separately or mark out-of-scope.
|
||||||
|
|
||||||
|
## API Tokens
|
||||||
|
- Files: `templates/token_create.gohtml`, `auth.APIKey` / `auth.CreatedAPIKey` types, `auth.Repository.CreateAPIKey` / `ValidateAPIKey` / `ListAPIKeys` / `RevokeAPIKey`
|
||||||
|
- Plan gap: not in M3. Belongs as an M3 addendum (developer access).
|
||||||
|
|
||||||
|
## Password Reset Flow
|
||||||
|
- Files: `templates/password_reset.gohtml`, `auth.PasswordResetToken`, `auth.Repository.CreatePasswordResetToken` / `CompletePasswordReset` / `ExpirePasswordResetToken`
|
||||||
|
- Plan gap: M3 task list does not mention password reset. Add as M3 addendum.
|
||||||
|
|
||||||
|
## Ops / Deployment Scaffolding
|
||||||
|
- Files: `entrypoint.sh`, `docker-compose.yml` updates (`c7dafae fix ports`, `81ceba1 port env`, `4666000 port fix`)
|
||||||
|
- Commit: `04a1f2b ops design and app`
|
||||||
|
- Plan gap: M6 operations. Early scaffolding only; full M6 work still pending.
|
||||||
@@ -24,12 +24,13 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
type App struct {
|
type App struct {
|
||||||
cfg config.Config
|
cfg config.Config
|
||||||
logger *slog.Logger
|
logger *slog.Logger
|
||||||
db database.DB
|
db database.DB
|
||||||
docs *docs.Service
|
docs *docs.Service
|
||||||
hub *realtime.Hub
|
hub *realtime.Hub
|
||||||
server *http.Server
|
server *http.Server
|
||||||
|
emailQueue *email.Queue
|
||||||
}
|
}
|
||||||
|
|
||||||
func New(ctx context.Context, cfg config.Config, logger *slog.Logger) (*App, error) {
|
func New(ctx context.Context, cfg config.Config, logger *slog.Logger) (*App, error) {
|
||||||
@@ -69,13 +70,33 @@ func New(ctx context.Context, cfg config.Config, logger *slog.Logger) (*App, err
|
|||||||
|
|
||||||
collabRepo := collaboration.NewRepository(db.SQL())
|
collabRepo := collaboration.NewRepository(db.SQL())
|
||||||
var emailSender collaboration.EmailSender
|
var emailSender collaboration.EmailSender
|
||||||
if cfg.Email.SMTPHost != "" {
|
var emailQueue *email.Queue
|
||||||
emailSender = email.NewSMTPSender(cfg.Email, logger)
|
emailRenderer := email.CollaborationRenderer{InstanceName: cfg.Email.InstanceName}
|
||||||
authService.SetEmailSender(emailSender)
|
switch cfg.Email.ResolvedProvider() {
|
||||||
} else {
|
case "postmark":
|
||||||
emailSender = email.NewNoOpSender(logger)
|
pm := email.NewPostmarkSender(cfg.Email.PostmarkToken, cfg.Email.From, cfg.Email.PostmarkAPIURL, logger)
|
||||||
|
emailSender = pm
|
||||||
|
authService.SetEmailSender(pm)
|
||||||
|
case "smtp":
|
||||||
|
sm := email.NewSMTPSender(cfg.Email, logger)
|
||||||
|
emailSender = sm
|
||||||
|
authService.SetEmailSender(sm)
|
||||||
|
default:
|
||||||
|
noOp := email.NewNoOpSender(logger)
|
||||||
|
emailSender = noOp
|
||||||
|
authService.SetEmailSender(noOp)
|
||||||
}
|
}
|
||||||
collabService := collaboration.NewService(collabRepo, emailSender, hub, logger)
|
// The durable queue wraps whichever sender is configured so retries
|
||||||
|
// survive process restarts. It is nil-safe in the service: if absent,
|
||||||
|
// notifications stay in-app only.
|
||||||
|
emailQueue = email.NewQueue(db.SQL(), emailSender, logger)
|
||||||
|
collabService := collaboration.NewServiceWithOptions(collabRepo, hub, logger, collaboration.Options{
|
||||||
|
EmailSender: emailSender,
|
||||||
|
EmailQueue: emailQueue,
|
||||||
|
EmailRenderer: emailRenderer,
|
||||||
|
UserLookup: authRepo,
|
||||||
|
PublicOrigin: cfg.Auth.PublicOrigin,
|
||||||
|
})
|
||||||
|
|
||||||
service.OnChange(func(change docs.DocumentChange) {
|
service.OnChange(func(change docs.DocumentChange) {
|
||||||
hub.Broadcast(realtime.Event{Type: "document_version", Data: change})
|
hub.Broadcast(realtime.Event{Type: "document_version", Data: change})
|
||||||
@@ -110,18 +131,22 @@ func New(ctx context.Context, cfg config.Config, logger *slog.Logger) (*App, err
|
|||||||
}
|
}
|
||||||
|
|
||||||
return &App{
|
return &App{
|
||||||
cfg: cfg,
|
cfg: cfg,
|
||||||
logger: logger,
|
logger: logger,
|
||||||
db: db,
|
db: db,
|
||||||
docs: service,
|
docs: service,
|
||||||
hub: hub,
|
hub: hub,
|
||||||
server: server,
|
server: server,
|
||||||
|
emailQueue: emailQueue,
|
||||||
}, nil
|
}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (a *App) Run(ctx context.Context) error {
|
func (a *App) Run(ctx context.Context) error {
|
||||||
go a.watchDocuments(ctx)
|
go a.watchDocuments(ctx)
|
||||||
go a.syncPoll(ctx)
|
go a.syncPoll(ctx)
|
||||||
|
if a.emailQueue != nil {
|
||||||
|
go a.emailQueue.Run(ctx)
|
||||||
|
}
|
||||||
|
|
||||||
go func() {
|
go func() {
|
||||||
<-ctx.Done()
|
<-ctx.Done()
|
||||||
|
|||||||
@@ -10,10 +10,46 @@ import (
|
|||||||
"github.com/tim/cairnquire/apps/server/internal/realtime"
|
"github.com/tim/cairnquire/apps/server/internal/realtime"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
// EmailSender mirrors email.Sender without taking a direct dependency on the
|
||||||
|
// email package (which would create a cycle once the queue imports from
|
||||||
|
// collaboration's repository types).
|
||||||
type EmailSender interface {
|
type EmailSender interface {
|
||||||
Send(ctx context.Context, to []string, subject, body string) error
|
Send(ctx context.Context, to []string, subject, body string) error
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// EmailQueue is the minimal contract the collaboration service needs to
|
||||||
|
// enqueue an outgoing email. The email.Queue type satisfies this.
|
||||||
|
type EmailQueue interface {
|
||||||
|
Enqueue(ctx context.Context, notificationID, to, subject, body string) (string, error)
|
||||||
|
}
|
||||||
|
|
||||||
|
// EmailRenderer renders a plain-text email for a notification kind. The
|
||||||
|
// email.Render function satisfies this.
|
||||||
|
type EmailRenderer interface {
|
||||||
|
Render(kind string, data EmailData) (subject, body string, err error)
|
||||||
|
}
|
||||||
|
|
||||||
|
// EmailData is the email-package TemplateData projected through the
|
||||||
|
// collaboration boundary to avoid an import cycle.
|
||||||
|
type EmailData struct {
|
||||||
|
ActorName string
|
||||||
|
ActorEmail string
|
||||||
|
DocumentPath string
|
||||||
|
DocumentTitle string
|
||||||
|
CommentBody string
|
||||||
|
MentionText string
|
||||||
|
ConflictDesc string
|
||||||
|
ViewURL string
|
||||||
|
UnsubURL string
|
||||||
|
InstanceName string
|
||||||
|
}
|
||||||
|
|
||||||
|
// UserLookup resolves a user id to an email address and display name. The
|
||||||
|
// auth.Repository satisfies this.
|
||||||
|
type UserLookup interface {
|
||||||
|
GetUserByID(ctx context.Context, userID string) (auth.User, error)
|
||||||
|
}
|
||||||
|
|
||||||
type NoOpEmailSender struct{}
|
type NoOpEmailSender struct{}
|
||||||
|
|
||||||
func (n *NoOpEmailSender) Send(ctx context.Context, to []string, subject, body string) error {
|
func (n *NoOpEmailSender) Send(ctx context.Context, to []string, subject, body string) error {
|
||||||
@@ -22,21 +58,45 @@ func (n *NoOpEmailSender) Send(ctx context.Context, to []string, subject, body s
|
|||||||
}
|
}
|
||||||
|
|
||||||
type Service struct {
|
type Service struct {
|
||||||
repo *Repository
|
repo *Repository
|
||||||
email EmailSender
|
email EmailSender
|
||||||
hub *realtime.Hub
|
queue EmailQueue
|
||||||
logger *slog.Logger
|
renderer EmailRenderer
|
||||||
|
users UserLookup
|
||||||
|
hub *realtime.Hub
|
||||||
|
logger *slog.Logger
|
||||||
|
publicOrigin string
|
||||||
|
}
|
||||||
|
|
||||||
|
type Options struct {
|
||||||
|
EmailSender EmailSender
|
||||||
|
EmailQueue EmailQueue
|
||||||
|
EmailRenderer EmailRenderer
|
||||||
|
UserLookup UserLookup
|
||||||
|
PublicOrigin string
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewService(repo *Repository, email EmailSender, hub *realtime.Hub, logger *slog.Logger) *Service {
|
func NewService(repo *Repository, email EmailSender, hub *realtime.Hub, logger *slog.Logger) *Service {
|
||||||
if email == nil {
|
return NewServiceWithOptions(repo, hub, logger, Options{EmailSender: email})
|
||||||
email = &NoOpEmailSender{}
|
}
|
||||||
|
|
||||||
|
// NewServiceWithOptions constructs the service with email queue + renderer +
|
||||||
|
// user lookup wiring. Callers that want email delivery should pass all of
|
||||||
|
// EmailQueue, EmailRenderer, UserLookup; otherwise notifications stay
|
||||||
|
// in-app only.
|
||||||
|
func NewServiceWithOptions(repo *Repository, hub *realtime.Hub, logger *slog.Logger, opts Options) *Service {
|
||||||
|
if opts.EmailSender == nil {
|
||||||
|
opts.EmailSender = &NoOpEmailSender{}
|
||||||
}
|
}
|
||||||
return &Service{
|
return &Service{
|
||||||
repo: repo,
|
repo: repo,
|
||||||
email: email,
|
email: opts.EmailSender,
|
||||||
hub: hub,
|
queue: opts.EmailQueue,
|
||||||
logger: logger,
|
renderer: opts.EmailRenderer,
|
||||||
|
users: opts.UserLookup,
|
||||||
|
hub: hub,
|
||||||
|
logger: logger,
|
||||||
|
publicOrigin: opts.PublicOrigin,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -224,6 +284,11 @@ func (s *Service) NotifyDocumentChanged(ctx context.Context, documentID string,
|
|||||||
s.logger.Warn("create notification", "error", err)
|
s.logger.Warn("create notification", "error", err)
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
s.enqueueNotificationEmail(ctx, notification.ID, w.UserID, "file_changed", EmailData{
|
||||||
|
ActorName: actorName,
|
||||||
|
DocumentPath: documentPath,
|
||||||
|
ViewURL: s.documentURL(documentPath),
|
||||||
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
// Broadcast bell update to all connected clients
|
// Broadcast bell update to all connected clients
|
||||||
@@ -257,11 +322,52 @@ func (s *Service) notifyWatchersOfComment(ctx context.Context, comment Comment)
|
|||||||
}
|
}
|
||||||
if err := s.repo.CreateNotification(ctx, notification); err != nil {
|
if err := s.repo.CreateNotification(ctx, notification); err != nil {
|
||||||
s.logger.Warn("create comment notification", "error", err)
|
s.logger.Warn("create comment notification", "error", err)
|
||||||
|
continue
|
||||||
}
|
}
|
||||||
|
s.enqueueNotificationEmail(ctx, notification.ID, w.UserID, "comment", EmailData{
|
||||||
|
ActorName: comment.AuthorName,
|
||||||
|
CommentBody: comment.Content,
|
||||||
|
ViewURL: s.documentURL(comment.DocumentID),
|
||||||
|
})
|
||||||
}
|
}
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// enqueueNotificationEmail resolves the watcher's email address, renders the
|
||||||
|
// template, and enqueues a durable email. Failures are logged and swallowed
|
||||||
|
// so a bad email config never blocks the in-app notification path.
|
||||||
|
func (s *Service) enqueueNotificationEmail(ctx context.Context, notificationID, userID, kind string, data EmailData) {
|
||||||
|
if s.queue == nil || s.renderer == nil || s.users == nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
user, err := s.users.GetUserByID(ctx, userID)
|
||||||
|
if err != nil {
|
||||||
|
s.logger.Warn("email: lookup user", "user_id", userID, "error", err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if user.Email == "" || user.Disabled {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if data.ActorName == "" {
|
||||||
|
data.ActorName = user.DisplayName
|
||||||
|
}
|
||||||
|
subject, body, err := s.renderer.Render(kind, data)
|
||||||
|
if err != nil {
|
||||||
|
s.logger.Warn("email: render", "kind", kind, "error", err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if _, err := s.queue.Enqueue(ctx, notificationID, user.Email, subject, body); err != nil {
|
||||||
|
s.logger.Warn("email: enqueue", "user_id", userID, "error", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *Service) documentURL(documentPath string) string {
|
||||||
|
if s.publicOrigin == "" {
|
||||||
|
return ""
|
||||||
|
}
|
||||||
|
return s.publicOrigin + "/" + documentPath
|
||||||
|
}
|
||||||
|
|
||||||
func (s *Service) broadcastNotification(userID string) {
|
func (s *Service) broadcastNotification(userID string) {
|
||||||
if s.hub == nil {
|
if s.hub == nil {
|
||||||
return
|
return
|
||||||
|
|||||||
@@ -38,9 +38,31 @@ type AuthConfig struct {
|
|||||||
}
|
}
|
||||||
|
|
||||||
type EmailConfig struct {
|
type EmailConfig struct {
|
||||||
SMTPHost string `json:"smtpHost"`
|
// Provider selects the sender implementation. Valid values: "noop",
|
||||||
SMTPPort int `json:"smtpPort"`
|
// "smtp", "postmark". When empty, the loader infers from the other
|
||||||
From string `json:"from"`
|
// fields: PostmarkToken wins, then SMTPHost, then NoOp.
|
||||||
|
Provider string `json:"provider"`
|
||||||
|
SMTPHost string `json:"smtpHost"`
|
||||||
|
SMTPPort int `json:"smtpPort"`
|
||||||
|
PostmarkToken string `json:"postmarkToken"`
|
||||||
|
PostmarkAPIURL string `json:"postmarkApiUrl"`
|
||||||
|
From string `json:"from"`
|
||||||
|
InstanceName string `json:"instanceName"`
|
||||||
|
}
|
||||||
|
|
||||||
|
// ResolvedProvider returns the concrete provider choice after applying the
|
||||||
|
// inference rule. Operators can set Provider explicitly to force a choice.
|
||||||
|
func (e EmailConfig) ResolvedProvider() string {
|
||||||
|
if e.Provider != "" {
|
||||||
|
return e.Provider
|
||||||
|
}
|
||||||
|
if e.PostmarkToken != "" {
|
||||||
|
return "postmark"
|
||||||
|
}
|
||||||
|
if e.SMTPHost != "" {
|
||||||
|
return "smtp"
|
||||||
|
}
|
||||||
|
return "noop"
|
||||||
}
|
}
|
||||||
|
|
||||||
func Default() Config {
|
func Default() Config {
|
||||||
@@ -59,9 +81,10 @@ func Default() Config {
|
|||||||
PublicOrigin: "http://localhost:8080",
|
PublicOrigin: "http://localhost:8080",
|
||||||
},
|
},
|
||||||
Email: EmailConfig{
|
Email: EmailConfig{
|
||||||
SMTPHost: "localhost",
|
SMTPHost: "localhost",
|
||||||
SMTPPort: 1025,
|
SMTPPort: 1025,
|
||||||
From: "notifications@cairnquire.local",
|
From: "notifications@cairnquire.local",
|
||||||
|
InstanceName: "Cairnquire",
|
||||||
},
|
},
|
||||||
LogLevel: "INFO",
|
LogLevel: "INFO",
|
||||||
}
|
}
|
||||||
@@ -87,6 +110,10 @@ func Load() (Config, error) {
|
|||||||
overrideString(&cfg.Content.StoreDir, "CAIRNQUIRE_CONTENT_STORE_DIR")
|
overrideString(&cfg.Content.StoreDir, "CAIRNQUIRE_CONTENT_STORE_DIR")
|
||||||
overrideString(&cfg.Auth.PublicOrigin, "CAIRNQUIRE_PUBLIC_ORIGIN")
|
overrideString(&cfg.Auth.PublicOrigin, "CAIRNQUIRE_PUBLIC_ORIGIN")
|
||||||
overrideString(&cfg.Email.SMTPHost, "CAIRNQUIRE_EMAIL_SMTP_HOST")
|
overrideString(&cfg.Email.SMTPHost, "CAIRNQUIRE_EMAIL_SMTP_HOST")
|
||||||
|
overrideString(&cfg.Email.PostmarkToken, "CAIRNQUIRE_EMAIL_POSTMARK_TOKEN")
|
||||||
|
overrideString(&cfg.Email.PostmarkAPIURL, "CAIRNQUIRE_EMAIL_POSTMARK_API_URL")
|
||||||
|
overrideString(&cfg.Email.Provider, "CAIRNQUIRE_EMAIL_PROVIDER")
|
||||||
|
overrideString(&cfg.Email.InstanceName, "CAIRNQUIRE_INSTANCE_NAME")
|
||||||
if port := os.Getenv("CAIRNQUIRE_EMAIL_SMTP_PORT"); port != "" {
|
if port := os.Getenv("CAIRNQUIRE_EMAIL_SMTP_PORT"); port != "" {
|
||||||
if p, err := strconv.Atoi(port); err == nil {
|
if p, err := strconv.Atoi(port); err == nil {
|
||||||
cfg.Email.SMTPPort = p
|
cfg.Email.SMTPPort = p
|
||||||
|
|||||||
@@ -0,0 +1,3 @@
|
|||||||
|
DROP INDEX IF EXISTS idx_email_queue_notification;
|
||||||
|
DROP INDEX IF EXISTS idx_email_queue_status;
|
||||||
|
DROP TABLE IF EXISTS email_queue;
|
||||||
@@ -0,0 +1,18 @@
|
|||||||
|
CREATE TABLE IF NOT EXISTS email_queue (
|
||||||
|
id TEXT PRIMARY KEY,
|
||||||
|
notification_id TEXT,
|
||||||
|
to_address TEXT NOT NULL,
|
||||||
|
subject TEXT NOT NULL,
|
||||||
|
body TEXT NOT NULL,
|
||||||
|
status TEXT NOT NULL DEFAULT 'pending' CHECK(status IN ('pending', 'sent', 'failed', 'skipped')),
|
||||||
|
attempts INTEGER NOT NULL DEFAULT 0,
|
||||||
|
max_attempts INTEGER NOT NULL DEFAULT 3,
|
||||||
|
last_error TEXT,
|
||||||
|
next_attempt_at TEXT NOT NULL,
|
||||||
|
sent_at TEXT,
|
||||||
|
created_at TEXT NOT NULL,
|
||||||
|
FOREIGN KEY(notification_id) REFERENCES notifications(id) ON DELETE SET NULL
|
||||||
|
);
|
||||||
|
|
||||||
|
CREATE INDEX IF NOT EXISTS idx_email_queue_status ON email_queue(status, next_attempt_at);
|
||||||
|
CREATE INDEX IF NOT EXISTS idx_email_queue_notification ON email_queue(notification_id);
|
||||||
29
apps/server/internal/email/collaboration_adapter.go
Normal file
29
apps/server/internal/email/collaboration_adapter.go
Normal file
@@ -0,0 +1,29 @@
|
|||||||
|
package email
|
||||||
|
|
||||||
|
import "github.com/tim/cairnquire/apps/server/internal/collaboration"
|
||||||
|
|
||||||
|
// CollaborationRenderer adapts email.Render to the collaboration.EmailRenderer
|
||||||
|
// interface by converting collaboration.EmailData into email.TemplateData.
|
||||||
|
type CollaborationRenderer struct {
|
||||||
|
InstanceName string
|
||||||
|
}
|
||||||
|
|
||||||
|
// Render satisfies collaboration.EmailRenderer.
|
||||||
|
func (c CollaborationRenderer) Render(kind string, data collaboration.EmailData) (string, string, error) {
|
||||||
|
td := TemplateData{
|
||||||
|
ActorName: data.ActorName,
|
||||||
|
ActorEmail: data.ActorEmail,
|
||||||
|
DocumentPath: data.DocumentPath,
|
||||||
|
DocumentTitle: data.DocumentTitle,
|
||||||
|
CommentBody: data.CommentBody,
|
||||||
|
MentionText: data.MentionText,
|
||||||
|
ConflictDesc: data.ConflictDesc,
|
||||||
|
ViewURL: data.ViewURL,
|
||||||
|
UnsubURL: data.UnsubURL,
|
||||||
|
InstanceName: c.InstanceName,
|
||||||
|
}
|
||||||
|
if td.InstanceName == "" {
|
||||||
|
td.InstanceName = "Cairnquire"
|
||||||
|
}
|
||||||
|
return Render(kind, td)
|
||||||
|
}
|
||||||
14
apps/server/internal/email/internal.go
Normal file
14
apps/server/internal/email/internal.go
Normal file
@@ -0,0 +1,14 @@
|
|||||||
|
package email
|
||||||
|
|
||||||
|
import (
|
||||||
|
"crypto/rand"
|
||||||
|
"encoding/base64"
|
||||||
|
)
|
||||||
|
|
||||||
|
func randomToken(length int) (string, error) {
|
||||||
|
buf := make([]byte, length)
|
||||||
|
if _, err := rand.Read(buf); err != nil {
|
||||||
|
return "", err
|
||||||
|
}
|
||||||
|
return base64.RawURLEncoding.EncodeToString(buf), nil
|
||||||
|
}
|
||||||
141
apps/server/internal/email/postmark.go
Normal file
141
apps/server/internal/email/postmark.go
Normal file
@@ -0,0 +1,141 @@
|
|||||||
|
package email
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bytes"
|
||||||
|
"context"
|
||||||
|
"encoding/json"
|
||||||
|
"fmt"
|
||||||
|
"io"
|
||||||
|
"log/slog"
|
||||||
|
"net/http"
|
||||||
|
"strings"
|
||||||
|
"time"
|
||||||
|
)
|
||||||
|
|
||||||
|
const (
|
||||||
|
defaultPostmarkURL = "https://api.postmarkapp.com"
|
||||||
|
postmarkTimeout = 15 * time.Second
|
||||||
|
)
|
||||||
|
|
||||||
|
// PostmarkSender sends email via the Postmark API.
|
||||||
|
type PostmarkSender struct {
|
||||||
|
serverToken string
|
||||||
|
apiURL string
|
||||||
|
from string
|
||||||
|
client *http.Client
|
||||||
|
logger *slog.Logger
|
||||||
|
}
|
||||||
|
|
||||||
|
// NewPostmarkSender constructs a Postmark sender. apiURL may be empty to use
|
||||||
|
// the public Postmark endpoint; tests can point it at an httptest server.
|
||||||
|
func NewPostmarkSender(serverToken, from, apiURL string, logger *slog.Logger) *PostmarkSender {
|
||||||
|
if apiURL == "" {
|
||||||
|
apiURL = defaultPostmarkURL
|
||||||
|
}
|
||||||
|
return &PostmarkSender{
|
||||||
|
serverToken: serverToken,
|
||||||
|
apiURL: strings.TrimRight(apiURL, "/"),
|
||||||
|
from: from,
|
||||||
|
client: &http.Client{Timeout: postmarkTimeout},
|
||||||
|
logger: logger,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
type postmarkRequest struct {
|
||||||
|
From string `json:"From"`
|
||||||
|
To string `json:"To"`
|
||||||
|
Subject string `json:"Subject"`
|
||||||
|
TextBody string `json:"TextBody"`
|
||||||
|
Headers []postmarkHeader `json:"Headers,omitempty"`
|
||||||
|
}
|
||||||
|
|
||||||
|
type postmarkHeader struct {
|
||||||
|
Name string `json:"Name"`
|
||||||
|
Value string `json:"Value"`
|
||||||
|
}
|
||||||
|
|
||||||
|
type postmarkResponse struct {
|
||||||
|
ErrorCode int `json:"ErrorCode"`
|
||||||
|
Message string `json:"Message"`
|
||||||
|
To string `json:"To"`
|
||||||
|
}
|
||||||
|
|
||||||
|
// WithReplyTo sets a Reply-To header on the outgoing message. Returns a copy
|
||||||
|
// of the sender with the reply-to address applied for the next Send call.
|
||||||
|
// Not safe for concurrent use; callers should construct a fresh sender per
|
||||||
|
// send if they need different reply-to addresses.
|
||||||
|
func (p *PostmarkSender) Send(ctx context.Context, to []string, subject, body string) error {
|
||||||
|
if p.serverToken == "" {
|
||||||
|
return fmt.Errorf("postmark server token not configured")
|
||||||
|
}
|
||||||
|
if p.from == "" {
|
||||||
|
return fmt.Errorf("postmark from address not configured")
|
||||||
|
}
|
||||||
|
if len(to) == 0 {
|
||||||
|
return fmt.Errorf("postmark send requires at least one recipient")
|
||||||
|
}
|
||||||
|
|
||||||
|
payload := postmarkRequest{
|
||||||
|
From: p.from,
|
||||||
|
To: strings.Join(to, ","),
|
||||||
|
Subject: subject,
|
||||||
|
TextBody: body,
|
||||||
|
}
|
||||||
|
bodyBytes, err := json.Marshal(payload)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("marshal postmark request: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
req, err := http.NewRequestWithContext(ctx, http.MethodPost, p.apiURL+"/email", bytes.NewReader(bodyBytes))
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("build postmark request: %w", err)
|
||||||
|
}
|
||||||
|
req.Header.Set("Accept", "application/json")
|
||||||
|
req.Header.Set("Content-Type", "application/json")
|
||||||
|
req.Header.Set("X-Postmark-Server-Token", p.serverToken)
|
||||||
|
|
||||||
|
resp, err := p.client.Do(req)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("postmark request: %w", err)
|
||||||
|
}
|
||||||
|
defer resp.Body.Close()
|
||||||
|
|
||||||
|
respBytes, _ := io.ReadAll(io.LimitReader(resp.Body, 1<<14))
|
||||||
|
|
||||||
|
if resp.StatusCode >= 500 {
|
||||||
|
return &TemporaryError{Cause: fmt.Errorf("postmark server error: status=%d body=%s", resp.StatusCode, truncate(string(respBytes)))}
|
||||||
|
}
|
||||||
|
|
||||||
|
if resp.StatusCode != http.StatusOK {
|
||||||
|
return fmt.Errorf("postmark rejected: status=%d body=%s", resp.StatusCode, truncate(string(respBytes)))
|
||||||
|
}
|
||||||
|
|
||||||
|
var parsed postmarkResponse
|
||||||
|
if err := json.Unmarshal(respBytes, &parsed); err != nil {
|
||||||
|
return fmt.Errorf("decode postmark response: %w", err)
|
||||||
|
}
|
||||||
|
if parsed.ErrorCode != 0 {
|
||||||
|
// Postmark distinguishes "inactive recipient" (406) and other errors. Treat
|
||||||
|
// 406 as non-retryable; others bubble up as generic failures.
|
||||||
|
return fmt.Errorf("postmark error code=%d: %s", parsed.ErrorCode, parsed.Message)
|
||||||
|
}
|
||||||
|
|
||||||
|
p.logger.Info("postmark email sent", "to", parsed.To, "subject", subject)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func truncate(s string) string {
|
||||||
|
if len(s) > 256 {
|
||||||
|
return s[:256] + "..."
|
||||||
|
}
|
||||||
|
return s
|
||||||
|
}
|
||||||
|
|
||||||
|
// TemporaryError wraps a failure that callers may retry. The queue worker
|
||||||
|
// uses errors.As to detect it and apply exponential backoff.
|
||||||
|
type TemporaryError struct {
|
||||||
|
Cause error
|
||||||
|
}
|
||||||
|
|
||||||
|
func (t *TemporaryError) Error() string { return t.Cause.Error() }
|
||||||
|
func (t *TemporaryError) Unwrap() error { return t.Cause }
|
||||||
113
apps/server/internal/email/postmark_test.go
Normal file
113
apps/server/internal/email/postmark_test.go
Normal file
@@ -0,0 +1,113 @@
|
|||||||
|
package email
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"encoding/json"
|
||||||
|
"errors"
|
||||||
|
"io"
|
||||||
|
"net/http"
|
||||||
|
"net/http/httptest"
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
)
|
||||||
|
|
||||||
|
func TestPostmarkSender_Send(t *testing.T) {
|
||||||
|
var gotBody postmarkRequest
|
||||||
|
var gotToken string
|
||||||
|
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||||
|
gotToken = r.Header.Get("X-Postmark-Server-Token")
|
||||||
|
body, _ := io.ReadAll(r.Body)
|
||||||
|
_ = json.Unmarshal(body, &gotBody)
|
||||||
|
w.Header().Set("Content-Type", "application/json")
|
||||||
|
_, _ = w.Write([]byte(`{"ErrorCode":0,"Message":"OK","To":"alice@example.com"}`))
|
||||||
|
}))
|
||||||
|
defer srv.Close()
|
||||||
|
|
||||||
|
sender := NewPostmarkSender("test-token", "noreply@example.com", srv.URL, testLogger(t))
|
||||||
|
if err := sender.Send(context.Background(), []string{"alice@example.com"}, "Hello", "World"); err != nil {
|
||||||
|
t.Fatalf("Send: %v", err)
|
||||||
|
}
|
||||||
|
if gotToken != "test-token" {
|
||||||
|
t.Errorf("token = %q, want test-token", gotToken)
|
||||||
|
}
|
||||||
|
if gotBody.From != "noreply@example.com" {
|
||||||
|
t.Errorf("from = %q", gotBody.From)
|
||||||
|
}
|
||||||
|
if gotBody.To != "alice@example.com" {
|
||||||
|
t.Errorf("to = %q", gotBody.To)
|
||||||
|
}
|
||||||
|
if gotBody.Subject != "Hello" {
|
||||||
|
t.Errorf("subject = %q", gotBody.Subject)
|
||||||
|
}
|
||||||
|
if gotBody.TextBody != "World" {
|
||||||
|
t.Errorf("body = %q", gotBody.TextBody)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestPostmarkSender_ErrorCode(t *testing.T) {
|
||||||
|
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||||
|
w.Header().Set("Content-Type", "application/json")
|
||||||
|
_, _ = w.Write([]byte(`{"ErrorCode":406,"Message":"Inactive recipient"}`))
|
||||||
|
}))
|
||||||
|
defer srv.Close()
|
||||||
|
|
||||||
|
sender := NewPostmarkSender("tok", "noreply@example.com", srv.URL, testLogger(t))
|
||||||
|
err := sender.Send(context.Background(), []string{"x@x.com"}, "s", "b")
|
||||||
|
if err == nil {
|
||||||
|
t.Fatal("expected error, got nil")
|
||||||
|
}
|
||||||
|
if !strings.Contains(err.Error(), "406") {
|
||||||
|
t.Errorf("error should mention code 406: %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestPostmarkSender_ServerError_Temporary(t *testing.T) {
|
||||||
|
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||||
|
w.WriteHeader(http.StatusInternalServerError)
|
||||||
|
_, _ = w.Write([]byte("boom"))
|
||||||
|
}))
|
||||||
|
defer srv.Close()
|
||||||
|
|
||||||
|
sender := NewPostmarkSender("tok", "noreply@example.com", srv.URL, testLogger(t))
|
||||||
|
err := sender.Send(context.Background(), []string{"x@x.com"}, "s", "b")
|
||||||
|
if err == nil {
|
||||||
|
t.Fatal("expected error")
|
||||||
|
}
|
||||||
|
var temp *TemporaryError
|
||||||
|
if !errors.As(err, &temp) {
|
||||||
|
t.Errorf("expected TemporaryError, got %T: %v", err, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestPostmarkSender_Validation(t *testing.T) {
|
||||||
|
sender := NewPostmarkSender("", "noreply@example.com", "", testLogger(t))
|
||||||
|
if err := sender.Send(context.Background(), []string{"x@x.com"}, "s", "b"); err == nil {
|
||||||
|
t.Fatal("expected error for missing token")
|
||||||
|
}
|
||||||
|
sender = NewPostmarkSender("tok", "", "", testLogger(t))
|
||||||
|
if err := sender.Send(context.Background(), []string{"x@x.com"}, "s", "b"); err == nil {
|
||||||
|
t.Fatal("expected error for missing from")
|
||||||
|
}
|
||||||
|
sender = NewPostmarkSender("tok", "noreply@example.com", "", testLogger(t))
|
||||||
|
if err := sender.Send(context.Background(), nil, "s", "b"); err == nil {
|
||||||
|
t.Fatal("expected error for empty recipients")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestPostmarkSender_MultipleRecipients(t *testing.T) {
|
||||||
|
var gotBody postmarkRequest
|
||||||
|
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||||
|
body, _ := io.ReadAll(r.Body)
|
||||||
|
_ = json.Unmarshal(body, &gotBody)
|
||||||
|
_, _ = w.Write([]byte(`{"ErrorCode":0,"Message":"OK"}`))
|
||||||
|
}))
|
||||||
|
defer srv.Close()
|
||||||
|
|
||||||
|
sender := NewPostmarkSender("tok", "noreply@example.com", srv.URL, testLogger(t))
|
||||||
|
if err := sender.Send(context.Background(), []string{"a@x.com", "b@x.com"}, "s", "b"); err != nil {
|
||||||
|
t.Fatalf("Send: %v", err)
|
||||||
|
}
|
||||||
|
if !strings.Contains(gotBody.To, "a@x.com") || !strings.Contains(gotBody.To, "b@x.com") {
|
||||||
|
t.Errorf("to = %q, expected both recipients", gotBody.To)
|
||||||
|
}
|
||||||
|
}
|
||||||
344
apps/server/internal/email/queue.go
Normal file
344
apps/server/internal/email/queue.go
Normal file
@@ -0,0 +1,344 @@
|
|||||||
|
package email
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"database/sql"
|
||||||
|
"fmt"
|
||||||
|
"log/slog"
|
||||||
|
"math"
|
||||||
|
"time"
|
||||||
|
)
|
||||||
|
|
||||||
|
// Queue stores outgoing emails durably and dispatches them via a Sender.
|
||||||
|
// It is safe for concurrent use; claim uses an atomic UPDATE ... WHERE
|
||||||
|
// status='pending' to lease rows without races.
|
||||||
|
type Queue struct {
|
||||||
|
db *sql.DB
|
||||||
|
sender Sender
|
||||||
|
logger *slog.Logger
|
||||||
|
|
||||||
|
pollInterval time.Duration
|
||||||
|
maxBackoff time.Duration
|
||||||
|
}
|
||||||
|
|
||||||
|
// QueueRow is the on-disk representation of a queued email.
|
||||||
|
type QueueRow struct {
|
||||||
|
ID string
|
||||||
|
NotificationID string
|
||||||
|
ToAddress string
|
||||||
|
Subject string
|
||||||
|
Body string
|
||||||
|
Status string
|
||||||
|
Attempts int
|
||||||
|
MaxAttempts int
|
||||||
|
LastError string
|
||||||
|
NextAttemptAt time.Time
|
||||||
|
SentAt *time.Time
|
||||||
|
CreatedAt time.Time
|
||||||
|
}
|
||||||
|
|
||||||
|
// NewQueue constructs a queue. sender may be nil; in that case emails stay
|
||||||
|
// pending forever (useful for tests that only inspect enqueue state).
|
||||||
|
func NewQueue(db *sql.DB, sender Sender, logger *slog.Logger) *Queue {
|
||||||
|
if logger == nil {
|
||||||
|
logger = slog.Default()
|
||||||
|
}
|
||||||
|
return &Queue{
|
||||||
|
db: db,
|
||||||
|
sender: sender,
|
||||||
|
logger: logger,
|
||||||
|
pollInterval: 10 * time.Second,
|
||||||
|
maxBackoff: 30 * time.Minute,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// SetPollInterval overrides the worker poll interval. Must be called before
|
||||||
|
// Run. Useful for tests.
|
||||||
|
func (q *Queue) SetPollInterval(d time.Duration) {
|
||||||
|
if d > 0 {
|
||||||
|
q.pollInterval = d
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// SetMaxBackoff overrides the maximum retry backoff.
|
||||||
|
func (q *Queue) SetMaxBackoff(d time.Duration) {
|
||||||
|
if d > 0 {
|
||||||
|
q.maxBackoff = d
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Enqueue inserts a new outgoing email. notificationID may be empty. The
|
||||||
|
// email is immediately eligible for delivery (next_attempt_at = now).
|
||||||
|
func (q *Queue) Enqueue(ctx context.Context, notificationID, to, subject, body string) (string, error) {
|
||||||
|
if to == "" {
|
||||||
|
return "", fmt.Errorf("enqueue: to address is required")
|
||||||
|
}
|
||||||
|
if subject == "" {
|
||||||
|
return "", fmt.Errorf("enqueue: subject is required")
|
||||||
|
}
|
||||||
|
id, err := randomQueueID()
|
||||||
|
if err != nil {
|
||||||
|
return "", err
|
||||||
|
}
|
||||||
|
now := time.Now().UTC()
|
||||||
|
if _, err := q.db.ExecContext(ctx, `
|
||||||
|
INSERT INTO email_queue (id, notification_id, to_address, subject, body, status, attempts, max_attempts, next_attempt_at, created_at)
|
||||||
|
VALUES (?, ?, ?, ?, ?, 'pending', 0, 3, ?, ?)
|
||||||
|
`, id, nullableString(notificationID), to, subject, body, formatTime(now), formatTime(now)); err != nil {
|
||||||
|
return "", fmt.Errorf("enqueue email: %w", err)
|
||||||
|
}
|
||||||
|
return id, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// MarkNotificationEmailed stamps notifications.emailed_at for the given
|
||||||
|
// notification id. Called by the worker once the email is dispatched.
|
||||||
|
func (q *Queue) MarkNotificationEmailed(ctx context.Context, notificationID string) error {
|
||||||
|
if notificationID == "" {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
_, err := q.db.ExecContext(ctx, `
|
||||||
|
UPDATE notifications SET emailed_at = ? WHERE id = ?
|
||||||
|
`, formatTime(time.Now().UTC()), notificationID)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("mark notification emailed: %w", err)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// Run starts the worker loop. It blocks until ctx is cancelled.
|
||||||
|
func (q *Queue) Run(ctx context.Context) {
|
||||||
|
q.logger.Info("email queue worker started", "poll_interval", q.pollInterval)
|
||||||
|
ticker := time.NewTicker(q.pollInterval)
|
||||||
|
defer ticker.Stop()
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case <-ctx.Done():
|
||||||
|
q.logger.Info("email queue worker stopped")
|
||||||
|
return
|
||||||
|
case <-ticker.C:
|
||||||
|
if err := q.processBatch(ctx); err != nil {
|
||||||
|
q.logger.Warn("email queue batch failed", "error", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// processBatch claims and sends up to 25 pending emails per tick.
|
||||||
|
func (q *Queue) processBatch(ctx context.Context) error {
|
||||||
|
if q.sender == nil {
|
||||||
|
// No sender configured; leave rows pending so a future sender can
|
||||||
|
// pick them up once wired.
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
rows, err := q.claimBatch(ctx, 25)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("claim batch: %w", err)
|
||||||
|
}
|
||||||
|
for _, row := range rows {
|
||||||
|
if err := q.deliver(ctx, row); err != nil {
|
||||||
|
q.logger.Warn("deliver email failed", "id", row.ID, "error", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// claimBatch atomically leases pending rows by flipping their status to
|
||||||
|
// 'pending' but bumping attempts so a concurrent worker cannot re-claim
|
||||||
|
// them until the next_attempt_at expires. We use a single UPDATE with a
|
||||||
|
// subquery to keep this race-free under SQLite.
|
||||||
|
func (q *Queue) claimBatch(ctx context.Context, limit int) ([]QueueRow, error) {
|
||||||
|
tx, err := q.db.BeginTx(ctx, nil)
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("begin claim: %w", err)
|
||||||
|
}
|
||||||
|
defer func() { _ = tx.Rollback() }()
|
||||||
|
|
||||||
|
rows, err := tx.QueryContext(ctx, `
|
||||||
|
SELECT id, COALESCE(notification_id, ''), to_address, subject, body, status, attempts, max_attempts,
|
||||||
|
COALESCE(last_error, ''), next_attempt_at, COALESCE(sent_at, ''), created_at
|
||||||
|
FROM email_queue
|
||||||
|
WHERE status = 'pending' AND next_attempt_at <= ?
|
||||||
|
ORDER BY next_attempt_at ASC
|
||||||
|
LIMIT ?
|
||||||
|
`, formatTime(time.Now().UTC()), limit)
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("query pending: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
var claimed []QueueRow
|
||||||
|
for rows.Next() {
|
||||||
|
row, err := scanQueueRow(rows)
|
||||||
|
if err != nil {
|
||||||
|
rows.Close()
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
claimed = append(claimed, row)
|
||||||
|
}
|
||||||
|
rows.Close()
|
||||||
|
if err := rows.Err(); err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
|
||||||
|
for _, row := range claimed {
|
||||||
|
if _, err := tx.ExecContext(ctx, `
|
||||||
|
UPDATE email_queue SET attempts = attempts + 1 WHERE id = ?
|
||||||
|
`, row.ID); err != nil {
|
||||||
|
return nil, fmt.Errorf("lease row %s: %w", row.ID, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if err := tx.Commit(); err != nil {
|
||||||
|
return nil, fmt.Errorf("commit claim: %w", err)
|
||||||
|
}
|
||||||
|
return claimed, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// deliver attempts one email. On success it marks the row sent and stamps
|
||||||
|
// the linked notification. On failure it schedules a retry with exponential
|
||||||
|
// backoff or marks the row failed if attempts are exhausted.
|
||||||
|
func (q *Queue) deliver(ctx context.Context, row QueueRow) error {
|
||||||
|
err := q.sender.Send(ctx, []string{row.ToAddress}, row.Subject, row.Body)
|
||||||
|
if err == nil {
|
||||||
|
if _, err := q.db.ExecContext(ctx, `
|
||||||
|
UPDATE email_queue SET status = 'sent', sent_at = ?, last_error = '' WHERE id = ?
|
||||||
|
`, formatTime(time.Now().UTC()), row.ID); err != nil {
|
||||||
|
return fmt.Errorf("mark sent: %w", err)
|
||||||
|
}
|
||||||
|
if err := q.MarkNotificationEmailed(ctx, row.NotificationID); err != nil {
|
||||||
|
q.logger.Warn("mark notification emailed", "id", row.NotificationID, "error", err)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// Failure: retry or give up.
|
||||||
|
if row.Attempts >= row.MaxAttempts {
|
||||||
|
if _, ferr := q.db.ExecContext(ctx, `
|
||||||
|
UPDATE email_queue SET status = 'failed', last_error = ? WHERE id = ?
|
||||||
|
`, truncateErr(err.Error()), row.ID); ferr != nil {
|
||||||
|
return fmt.Errorf("mark failed: %w", ferr)
|
||||||
|
}
|
||||||
|
q.logger.Warn("email permanently failed", "id", row.ID, "attempts", row.Attempts, "error", err)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
backoff := q.backoff(row.Attempts)
|
||||||
|
next := time.Now().UTC().Add(backoff)
|
||||||
|
if _, err := q.db.ExecContext(ctx, `
|
||||||
|
UPDATE email_queue SET status = 'pending', next_attempt_at = ?, last_error = ? WHERE id = ?
|
||||||
|
`, formatTime(next), truncateErr(err.Error()), row.ID); err != nil {
|
||||||
|
return fmt.Errorf("schedule retry: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// TemporaryError (Postmark 5xx) is retried with the same backoff as
|
||||||
|
// other errors; the max-attempts cap still bounds total work.
|
||||||
|
q.logger.Info("email retry scheduled", "id", row.ID, "attempt", row.Attempts, "backoff", backoff, "error", err)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (q *Queue) backoff(attempts int) time.Duration {
|
||||||
|
// 2^attempts seconds, capped: 2s, 4s, 8s, 16s, ... up to maxBackoff.
|
||||||
|
d := time.Duration(math.Pow(2, float64(attempts))) * time.Second
|
||||||
|
if d > q.maxBackoff {
|
||||||
|
d = q.maxBackoff
|
||||||
|
}
|
||||||
|
if d < time.Second {
|
||||||
|
d = time.Second
|
||||||
|
}
|
||||||
|
return d
|
||||||
|
}
|
||||||
|
|
||||||
|
// PendingCount returns the number of rows in pending status, primarily for
|
||||||
|
// tests and operator dashboards.
|
||||||
|
func (q *Queue) PendingCount(ctx context.Context) (int, error) {
|
||||||
|
var count int
|
||||||
|
err := q.db.QueryRowContext(ctx, `SELECT COUNT(1) FROM email_queue WHERE status = 'pending'`).Scan(&count)
|
||||||
|
if err != nil {
|
||||||
|
return 0, err
|
||||||
|
}
|
||||||
|
return count, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// ListFailed returns rows that have exhausted their retries, for operator
|
||||||
|
// inspection.
|
||||||
|
func (q *Queue) ListFailed(ctx context.Context, limit int) ([]QueueRow, error) {
|
||||||
|
if limit <= 0 {
|
||||||
|
limit = 50
|
||||||
|
}
|
||||||
|
rows, err := q.db.QueryContext(ctx, `
|
||||||
|
SELECT id, COALESCE(notification_id, ''), to_address, subject, body, status, attempts, max_attempts,
|
||||||
|
COALESCE(last_error, ''), next_attempt_at, COALESCE(sent_at, ''), created_at
|
||||||
|
FROM email_queue
|
||||||
|
WHERE status = 'failed'
|
||||||
|
ORDER BY created_at DESC
|
||||||
|
LIMIT ?
|
||||||
|
`, limit)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
defer rows.Close()
|
||||||
|
|
||||||
|
var out []QueueRow
|
||||||
|
for rows.Next() {
|
||||||
|
row, err := scanQueueRow(rows)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
out = append(out, row)
|
||||||
|
}
|
||||||
|
return out, rows.Err()
|
||||||
|
}
|
||||||
|
|
||||||
|
func scanQueueRow(rows *sql.Rows) (QueueRow, error) {
|
||||||
|
var r QueueRow
|
||||||
|
var sentAt, lastError string
|
||||||
|
var nextAttempt, createdAt string
|
||||||
|
if err := rows.Scan(&r.ID, &r.NotificationID, &r.ToAddress, &r.Subject, &r.Body, &r.Status,
|
||||||
|
&r.Attempts, &r.MaxAttempts, &lastError, &nextAttempt, &sentAt, &createdAt); err != nil {
|
||||||
|
return QueueRow{}, fmt.Errorf("scan queue row: %w", err)
|
||||||
|
}
|
||||||
|
r.LastError = lastError
|
||||||
|
var err error
|
||||||
|
r.NextAttemptAt, err = time.Parse(time.RFC3339, nextAttempt)
|
||||||
|
if err != nil {
|
||||||
|
return QueueRow{}, fmt.Errorf("parse next_attempt_at: %w", err)
|
||||||
|
}
|
||||||
|
r.CreatedAt, err = time.Parse(time.RFC3339, createdAt)
|
||||||
|
if err != nil {
|
||||||
|
return QueueRow{}, fmt.Errorf("parse created_at: %w", err)
|
||||||
|
}
|
||||||
|
if sentAt != "" {
|
||||||
|
t, err := time.Parse(time.RFC3339, sentAt)
|
||||||
|
if err == nil {
|
||||||
|
r.SentAt = &t
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return r, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func formatTime(t time.Time) string {
|
||||||
|
return t.UTC().Format(time.RFC3339)
|
||||||
|
}
|
||||||
|
|
||||||
|
func nullableString(s string) any {
|
||||||
|
if s == "" {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
return s
|
||||||
|
}
|
||||||
|
|
||||||
|
func truncateErr(s string) string {
|
||||||
|
if len(s) > 512 {
|
||||||
|
return s[:512] + "..."
|
||||||
|
}
|
||||||
|
return s
|
||||||
|
}
|
||||||
|
|
||||||
|
// randomQueueID produces an id with enough entropy to avoid collisions. We
|
||||||
|
// keep a small helper in-package to avoid an import cycle on auth.
|
||||||
|
func randomQueueID() (string, error) {
|
||||||
|
id, err := randomToken(18)
|
||||||
|
if err != nil {
|
||||||
|
return "", err
|
||||||
|
}
|
||||||
|
return "email:" + id, nil
|
||||||
|
}
|
||||||
248
apps/server/internal/email/queue_test.go
Normal file
248
apps/server/internal/email/queue_test.go
Normal file
@@ -0,0 +1,248 @@
|
|||||||
|
package email
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"database/sql"
|
||||||
|
"errors"
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/tim/cairnquire/apps/server/internal/config"
|
||||||
|
"github.com/tim/cairnquire/apps/server/internal/database"
|
||||||
|
)
|
||||||
|
|
||||||
|
func newTestDB(t *testing.T) *sql.DB {
|
||||||
|
t.Helper()
|
||||||
|
cfg := config.DatabaseConfig{Path: t.TempDir() + "/test.sqlite"}
|
||||||
|
handle, err := database.Open(context.Background(), cfg)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("open db: %v", err)
|
||||||
|
}
|
||||||
|
t.Cleanup(func() { _ = handle.Close() })
|
||||||
|
db := handle.SQL()
|
||||||
|
if err := database.ApplyMigrations(context.Background(), db); err != nil {
|
||||||
|
t.Fatalf("migrate: %v", err)
|
||||||
|
}
|
||||||
|
// Seed a user + notification so the queue's notification_id FK is
|
||||||
|
// satisfiable and MarkNotificationEmailed can update a real row.
|
||||||
|
_, err = db.Exec(`INSERT INTO users (id, email, display_name, role, created_at, last_seen_at) VALUES ('user:alice@example.com', 'alice@example.com', 'Alice', 'admin', ?, ?)`,
|
||||||
|
time.Now().UTC().Format(time.RFC3339), time.Now().UTC().Format(time.RFC3339))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("seed user: %v", err)
|
||||||
|
}
|
||||||
|
_, err = db.Exec(`INSERT INTO notifications (id, user_id, type, resource_type, resource_id, message, created_at) VALUES (?, 'user:alice@example.com', 'file_changed', 'document', 'doc:test', 'test', ?)`,
|
||||||
|
"notif:test-1", time.Now().UTC().Format(time.RFC3339))
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("seed notification: %v", err)
|
||||||
|
}
|
||||||
|
return db
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestQueue_EnqueueAndPendingCount(t *testing.T) {
|
||||||
|
db := newTestDB(t)
|
||||||
|
q := NewQueue(db, nil, testLogger(t))
|
||||||
|
|
||||||
|
id, err := q.Enqueue(context.Background(), "notif:test-1", "alice@example.com", "Subject", "Body")
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("Enqueue: %v", err)
|
||||||
|
}
|
||||||
|
if id == "" {
|
||||||
|
t.Fatal("expected non-empty id")
|
||||||
|
}
|
||||||
|
count, err := q.PendingCount(context.Background())
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("PendingCount: %v", err)
|
||||||
|
}
|
||||||
|
if count != 1 {
|
||||||
|
t.Errorf("pending = %d, want 1", count)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestQueue_EnqueueValidation(t *testing.T) {
|
||||||
|
db := newTestDB(t)
|
||||||
|
q := NewQueue(db, nil, testLogger(t))
|
||||||
|
|
||||||
|
// Empty notification_id is allowed (the column is nullable).
|
||||||
|
if _, err := q.Enqueue(context.Background(), "", "alice@example.com", "s", "b"); err != nil {
|
||||||
|
t.Fatalf("empty notification_id should be allowed: %v", err)
|
||||||
|
}
|
||||||
|
if _, err := q.Enqueue(context.Background(), "notif:test-1", "", "s", "b"); err == nil {
|
||||||
|
t.Fatal("expected error for empty to")
|
||||||
|
}
|
||||||
|
if _, err := q.Enqueue(context.Background(), "notif:test-1", "alice@example.com", "", "b"); err == nil {
|
||||||
|
t.Fatal("expected error for empty subject")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
type recordingSender struct {
|
||||||
|
sends []recordedSend
|
||||||
|
err error
|
||||||
|
}
|
||||||
|
|
||||||
|
type recordedSend struct {
|
||||||
|
to []string
|
||||||
|
subject string
|
||||||
|
body string
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *recordingSender) Send(ctx context.Context, to []string, subject, body string) error {
|
||||||
|
if r.err != nil {
|
||||||
|
return r.err
|
||||||
|
}
|
||||||
|
r.sends = append(r.sends, recordedSend{to: to, subject: subject, body: body})
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestQueue_DeliverSuccess(t *testing.T) {
|
||||||
|
db := newTestDB(t)
|
||||||
|
sender := &recordingSender{}
|
||||||
|
q := NewQueue(db, sender, testLogger(t))
|
||||||
|
q.SetPollInterval(20 * time.Millisecond)
|
||||||
|
|
||||||
|
if _, err := q.Enqueue(context.Background(), "notif:test-1", "alice@example.com", "Hello", "World"); err != nil {
|
||||||
|
t.Fatalf("Enqueue: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
|
||||||
|
defer cancel()
|
||||||
|
go q.Run(ctx)
|
||||||
|
|
||||||
|
if err := waitFor(ctx, func() bool {
|
||||||
|
n, _ := q.PendingCount(ctx)
|
||||||
|
return n == 0
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatalf("email never sent: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if len(sender.sends) != 1 {
|
||||||
|
t.Fatalf("sends = %d, want 1", len(sender.sends))
|
||||||
|
}
|
||||||
|
if sender.sends[0].subject != "Hello" {
|
||||||
|
t.Errorf("subject = %q", sender.sends[0].subject)
|
||||||
|
}
|
||||||
|
if sender.sends[0].body != "World" {
|
||||||
|
t.Errorf("body = %q", sender.sends[0].body)
|
||||||
|
}
|
||||||
|
|
||||||
|
// The linked notification should now have emailed_at stamped.
|
||||||
|
var emailedAt string
|
||||||
|
if err := db.QueryRow(`SELECT COALESCE(emailed_at, '') FROM notifications WHERE id = ?`, "notif:test-1").Scan(&emailedAt); err != nil {
|
||||||
|
t.Fatalf("query emailed_at: %v", err)
|
||||||
|
}
|
||||||
|
if emailedAt == "" {
|
||||||
|
t.Error("expected notifications.emailed_at to be set after send")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestQueue_RetryThenSucceed(t *testing.T) {
|
||||||
|
db := newTestDB(t)
|
||||||
|
// Fail twice, then succeed.
|
||||||
|
sender := &flakySender{failTimes: 2}
|
||||||
|
q := NewQueue(db, sender, testLogger(t))
|
||||||
|
q.SetPollInterval(15 * time.Millisecond)
|
||||||
|
q.SetMaxBackoff(50 * time.Millisecond)
|
||||||
|
|
||||||
|
if _, err := q.Enqueue(context.Background(), "notif:test-1", "alice@example.com", "Retry", "Body"); err != nil {
|
||||||
|
t.Fatalf("Enqueue: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
||||||
|
defer cancel()
|
||||||
|
go q.Run(ctx)
|
||||||
|
|
||||||
|
if err := waitFor(ctx, func() bool {
|
||||||
|
return sender.successCount >= 1
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatalf("email never succeeded: successCount=%d attempts=%d: %v", sender.successCount, sender.attempts, err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Verify the row is marked sent.
|
||||||
|
var status string
|
||||||
|
if err := db.QueryRow(`SELECT status FROM email_queue WHERE to_address = ?`, "alice@example.com").Scan(&status); err != nil {
|
||||||
|
t.Fatalf("query status: %v", err)
|
||||||
|
}
|
||||||
|
if status != "sent" {
|
||||||
|
t.Errorf("status = %q, want sent", status)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestQueue_PermanentFailure(t *testing.T) {
|
||||||
|
db := newTestDB(t)
|
||||||
|
sender := &recordingSender{err: errors.New("permanent Postmark rejection")}
|
||||||
|
q := NewQueue(db, sender, testLogger(t))
|
||||||
|
q.SetPollInterval(15 * time.Millisecond)
|
||||||
|
q.SetMaxBackoff(20 * time.Millisecond)
|
||||||
|
|
||||||
|
if _, err := q.Enqueue(context.Background(), "notif:test-1", "alice@example.com", "Fail", "Body"); err != nil {
|
||||||
|
t.Fatalf("Enqueue: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
||||||
|
defer cancel()
|
||||||
|
go q.Run(ctx)
|
||||||
|
|
||||||
|
if err := waitFor(ctx, func() bool {
|
||||||
|
failed, _ := q.ListFailed(ctx, 10)
|
||||||
|
return len(failed) == 1
|
||||||
|
}); err != nil {
|
||||||
|
failed, _ := q.ListFailed(ctx, 10)
|
||||||
|
t.Fatalf("expected 1 failed row, got %d: %v", len(failed), err)
|
||||||
|
}
|
||||||
|
|
||||||
|
failed, _ := q.ListFailed(ctx, 10)
|
||||||
|
if !strings.Contains(failed[0].LastError, "permanent") {
|
||||||
|
t.Errorf("last_error = %q", failed[0].LastError)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestQueue_NilSenderLeavesPending(t *testing.T) {
|
||||||
|
db := newTestDB(t)
|
||||||
|
q := NewQueue(db, nil, testLogger(t))
|
||||||
|
q.SetPollInterval(15 * time.Millisecond)
|
||||||
|
|
||||||
|
if _, err := q.Enqueue(context.Background(), "notif:test-1", "alice@example.com", "Pending", "Body"); err != nil {
|
||||||
|
t.Fatalf("Enqueue: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
ctx, cancel := context.WithTimeout(context.Background(), 300*time.Millisecond)
|
||||||
|
defer cancel()
|
||||||
|
go q.Run(ctx)
|
||||||
|
|
||||||
|
<-ctx.Done()
|
||||||
|
count, _ := q.PendingCount(context.Background())
|
||||||
|
if count != 1 {
|
||||||
|
t.Errorf("expected row to stay pending with nil sender, got %d", count)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// flakySender fails the first N attempts then succeeds.
|
||||||
|
type flakySender struct {
|
||||||
|
failTimes int
|
||||||
|
attempts int
|
||||||
|
successCount int
|
||||||
|
}
|
||||||
|
|
||||||
|
func (f *flakySender) Send(ctx context.Context, to []string, subject, body string) error {
|
||||||
|
f.attempts++
|
||||||
|
if f.attempts <= f.failTimes {
|
||||||
|
return errors.New("transient failure")
|
||||||
|
}
|
||||||
|
f.successCount++
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func waitFor(ctx context.Context, cond func() bool) error {
|
||||||
|
ticker := time.NewTicker(10 * time.Millisecond)
|
||||||
|
defer ticker.Stop()
|
||||||
|
for {
|
||||||
|
if cond() {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
select {
|
||||||
|
case <-ctx.Done():
|
||||||
|
return ctx.Err()
|
||||||
|
case <-ticker.C:
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
149
apps/server/internal/email/templates.go
Normal file
149
apps/server/internal/email/templates.go
Normal file
@@ -0,0 +1,149 @@
|
|||||||
|
package email
|
||||||
|
|
||||||
|
import (
|
||||||
|
"fmt"
|
||||||
|
"strings"
|
||||||
|
"time"
|
||||||
|
)
|
||||||
|
|
||||||
|
// TemplateData carries the context required to render any notification email.
|
||||||
|
// Fields are optional per template; callers should populate what they have.
|
||||||
|
type TemplateData struct {
|
||||||
|
ActorName string
|
||||||
|
ActorEmail string
|
||||||
|
DocumentPath string
|
||||||
|
DocumentTitle string
|
||||||
|
CommentBody string
|
||||||
|
MentionText string
|
||||||
|
ConflictDesc string
|
||||||
|
ViewURL string
|
||||||
|
UnsubURL string
|
||||||
|
InstanceName string
|
||||||
|
}
|
||||||
|
|
||||||
|
// Render returns a plain-text subject and body for the given notification
|
||||||
|
// type. Subjects are short and prefixed with the document path so email
|
||||||
|
// clients thread per-document.
|
||||||
|
func Render(kind string, data TemplateData) (subject, body string, err error) {
|
||||||
|
if data.InstanceName == "" {
|
||||||
|
data.InstanceName = "Cairnquire"
|
||||||
|
}
|
||||||
|
switch kind {
|
||||||
|
case "file_changed":
|
||||||
|
return renderFileChanged(data)
|
||||||
|
case "comment":
|
||||||
|
return renderComment(data)
|
||||||
|
case "mention":
|
||||||
|
return renderMention(data)
|
||||||
|
case "conflict":
|
||||||
|
return renderConflict(data)
|
||||||
|
default:
|
||||||
|
return "", "", fmt.Errorf("unknown email template kind %q", kind)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func subjectPrefix(data TemplateData) string {
|
||||||
|
if data.DocumentPath != "" {
|
||||||
|
return "[" + data.DocumentPath + "]"
|
||||||
|
}
|
||||||
|
return "[" + data.InstanceName + "]"
|
||||||
|
}
|
||||||
|
|
||||||
|
func footer(data TemplateData) string {
|
||||||
|
var b strings.Builder
|
||||||
|
b.WriteString("\n--\n")
|
||||||
|
if data.ViewURL != "" {
|
||||||
|
fmt.Fprintf(&b, "View online: %s\n", data.ViewURL)
|
||||||
|
}
|
||||||
|
if data.UnsubURL != "" {
|
||||||
|
fmt.Fprintf(&b, "Unsubscribe: %s\n", data.UnsubURL)
|
||||||
|
}
|
||||||
|
fmt.Fprintf(&b, "You receive this because you watch this document on %s.\n", data.InstanceName)
|
||||||
|
return b.String()
|
||||||
|
}
|
||||||
|
|
||||||
|
func renderFileChanged(data TemplateData) (string, string, error) {
|
||||||
|
title := data.DocumentTitle
|
||||||
|
if title == "" {
|
||||||
|
title = data.DocumentPath
|
||||||
|
}
|
||||||
|
actor := data.ActorName
|
||||||
|
if actor == "" {
|
||||||
|
actor = "Someone"
|
||||||
|
}
|
||||||
|
subject := fmt.Sprintf("%s Updated by %s", subjectPrefix(data), actor)
|
||||||
|
|
||||||
|
var b strings.Builder
|
||||||
|
fmt.Fprintf(&b, "%s updated %q at %s UTC.\n", actor, title, time.Now().UTC().Format("2006-01-02 15:04"))
|
||||||
|
if data.DocumentPath != "" {
|
||||||
|
fmt.Fprintf(&b, "Path: %s\n", data.DocumentPath)
|
||||||
|
}
|
||||||
|
b.WriteString("\nView the changes online to see what is new.\n")
|
||||||
|
b.WriteString(footer(data))
|
||||||
|
return subject, b.String(), nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func renderComment(data TemplateData) (string, string, error) {
|
||||||
|
actor := data.ActorName
|
||||||
|
if actor == "" {
|
||||||
|
actor = "Someone"
|
||||||
|
}
|
||||||
|
subject := fmt.Sprintf("%s Comment from %s", subjectPrefix(data), actor)
|
||||||
|
|
||||||
|
var b strings.Builder
|
||||||
|
fmt.Fprintf(&b, "%s commented on %q:\n", actor, docLabel(data))
|
||||||
|
b.WriteString("\n")
|
||||||
|
if data.CommentBody != "" {
|
||||||
|
// Quote the comment body, line by line, as plain text.
|
||||||
|
for _, line := range strings.Split(data.CommentBody, "\n") {
|
||||||
|
fmt.Fprintf(&b, "> %s\n", line)
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
b.WriteString("> (no body)\n")
|
||||||
|
}
|
||||||
|
b.WriteString("\nReply to this email to respond.\n")
|
||||||
|
b.WriteString(footer(data))
|
||||||
|
return subject, b.String(), nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func renderMention(data TemplateData) (string, string, error) {
|
||||||
|
actor := data.ActorName
|
||||||
|
if actor == "" {
|
||||||
|
actor = "Someone"
|
||||||
|
}
|
||||||
|
subject := fmt.Sprintf("%s Mention from %s", subjectPrefix(data), actor)
|
||||||
|
|
||||||
|
var b strings.Builder
|
||||||
|
fmt.Fprintf(&b, "%s mentioned you on %q.\n", actor, docLabel(data))
|
||||||
|
if data.MentionText != "" {
|
||||||
|
b.WriteString("\n")
|
||||||
|
for _, line := range strings.Split(data.MentionText, "\n") {
|
||||||
|
fmt.Fprintf(&b, "> %s\n", line)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
b.WriteString(footer(data))
|
||||||
|
return subject, b.String(), nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func renderConflict(data TemplateData) (string, string, error) {
|
||||||
|
subject := fmt.Sprintf("%s Sync conflict", subjectPrefix(data))
|
||||||
|
|
||||||
|
var b strings.Builder
|
||||||
|
fmt.Fprintf(&b, "A sync conflict was detected on %q.\n", docLabel(data))
|
||||||
|
if data.ConflictDesc != "" {
|
||||||
|
fmt.Fprintf(&b, "\nDetails: %s\n", data.ConflictDesc)
|
||||||
|
}
|
||||||
|
b.WriteString("\nOpen the document to review and resolve the conflict.\n")
|
||||||
|
b.WriteString(footer(data))
|
||||||
|
return subject, b.String(), nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func docLabel(data TemplateData) string {
|
||||||
|
if data.DocumentTitle != "" {
|
||||||
|
return data.DocumentTitle
|
||||||
|
}
|
||||||
|
if data.DocumentPath != "" {
|
||||||
|
return data.DocumentPath
|
||||||
|
}
|
||||||
|
return "a document"
|
||||||
|
}
|
||||||
119
apps/server/internal/email/templates_test.go
Normal file
119
apps/server/internal/email/templates_test.go
Normal file
@@ -0,0 +1,119 @@
|
|||||||
|
package email
|
||||||
|
|
||||||
|
import (
|
||||||
|
"strings"
|
||||||
|
"testing"
|
||||||
|
)
|
||||||
|
|
||||||
|
func TestRender_FileChanged(t *testing.T) {
|
||||||
|
subject, body, err := Render("file_changed", TemplateData{
|
||||||
|
ActorName: "Alice",
|
||||||
|
DocumentPath: "docs/api.md",
|
||||||
|
DocumentTitle: "API Reference",
|
||||||
|
ViewURL: "https://example.com/docs/api.md",
|
||||||
|
UnsubURL: "https://example.com/unsub/123",
|
||||||
|
InstanceName: "Cairnquire",
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("Render: %v", err)
|
||||||
|
}
|
||||||
|
if !strings.Contains(subject, "docs/api.md") {
|
||||||
|
t.Errorf("subject should include path: %q", subject)
|
||||||
|
}
|
||||||
|
if !strings.Contains(subject, "Alice") {
|
||||||
|
t.Errorf("subject should include actor: %q", subject)
|
||||||
|
}
|
||||||
|
if !strings.Contains(body, "Alice") {
|
||||||
|
t.Errorf("body should mention actor: %q", body)
|
||||||
|
}
|
||||||
|
if !strings.Contains(body, "API Reference") {
|
||||||
|
t.Errorf("body should mention title: %q", body)
|
||||||
|
}
|
||||||
|
if !strings.Contains(body, "https://example.com/docs/api.md") {
|
||||||
|
t.Errorf("body should include view url: %q", body)
|
||||||
|
}
|
||||||
|
if !strings.Contains(body, "https://example.com/unsub/123") {
|
||||||
|
t.Errorf("body should include unsub url: %q", body)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestRender_Comment(t *testing.T) {
|
||||||
|
subject, body, err := Render("comment", TemplateData{
|
||||||
|
ActorName: "Bob",
|
||||||
|
DocumentTitle: "Getting Started",
|
||||||
|
CommentBody: "This step is unclear.\nCan we add an example?",
|
||||||
|
ViewURL: "https://example.com/docs/getting-started",
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("Render: %v", err)
|
||||||
|
}
|
||||||
|
if !strings.Contains(subject, "Bob") {
|
||||||
|
t.Errorf("subject should include actor: %q", subject)
|
||||||
|
}
|
||||||
|
if !strings.Contains(body, "> This step is unclear.") {
|
||||||
|
t.Errorf("body should quote comment: %q", body)
|
||||||
|
}
|
||||||
|
if !strings.Contains(body, "> Can we add an example?") {
|
||||||
|
t.Errorf("body should quote second line: %q", body)
|
||||||
|
}
|
||||||
|
if !strings.Contains(body, "Reply to this email") {
|
||||||
|
t.Errorf("body should mention reply-to-respond: %q", body)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestRender_Mention(t *testing.T) {
|
||||||
|
subject, body, err := Render("mention", TemplateData{
|
||||||
|
ActorName: "Carol",
|
||||||
|
DocumentTitle: "Spec",
|
||||||
|
MentionText: "@dan please review",
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("Render: %v", err)
|
||||||
|
}
|
||||||
|
if !strings.Contains(subject, "Mention") {
|
||||||
|
t.Errorf("subject: %q", subject)
|
||||||
|
}
|
||||||
|
if !strings.Contains(body, "> @dan please review") {
|
||||||
|
t.Errorf("body should quote mention: %q", body)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestRender_Conflict(t *testing.T) {
|
||||||
|
subject, body, err := Render("conflict", TemplateData{
|
||||||
|
DocumentPath: "notes/2024.md",
|
||||||
|
ConflictDesc: "Local and remote both edited line 12",
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("Render: %v", err)
|
||||||
|
}
|
||||||
|
if !strings.Contains(subject, "conflict") {
|
||||||
|
t.Errorf("subject: %q", subject)
|
||||||
|
}
|
||||||
|
if !strings.Contains(body, "sync conflict") {
|
||||||
|
t.Errorf("body should mention conflict: %q", body)
|
||||||
|
}
|
||||||
|
if !strings.Contains(body, "line 12") {
|
||||||
|
t.Errorf("body should include conflict desc: %q", body)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestRender_UnknownKind(t *testing.T) {
|
||||||
|
_, _, err := Render("bogus", TemplateData{})
|
||||||
|
if err == nil {
|
||||||
|
t.Fatal("expected error for unknown kind")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestRender_PlainTextOnly(t *testing.T) {
|
||||||
|
for _, kind := range []string{"file_changed", "comment", "mention", "conflict"} {
|
||||||
|
_, body, err := Render(kind, TemplateData{
|
||||||
|
ActorName: "A", DocumentTitle: "T", CommentBody: "<b>html</b>",
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("Render %s: %v", kind, err)
|
||||||
|
}
|
||||||
|
if strings.Contains(body, "<html>") || strings.Contains(body, "<body>") {
|
||||||
|
t.Errorf("template %s produced HTML: %q", kind, body)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
18
apps/server/internal/email/test_helpers_test.go
Normal file
18
apps/server/internal/email/test_helpers_test.go
Normal file
@@ -0,0 +1,18 @@
|
|||||||
|
package email
|
||||||
|
|
||||||
|
import (
|
||||||
|
"log/slog"
|
||||||
|
"testing"
|
||||||
|
)
|
||||||
|
|
||||||
|
func testLogger(t *testing.T) *slog.Logger {
|
||||||
|
t.Helper()
|
||||||
|
return slog.New(slog.NewTextHandler(&testWriter{t: t}, &slog.HandlerOptions{Level: slog.LevelWarn}))
|
||||||
|
}
|
||||||
|
|
||||||
|
type testWriter struct{ t *testing.T }
|
||||||
|
|
||||||
|
func (w *testWriter) Write(p []byte) (int, error) {
|
||||||
|
w.t.Logf("%s", p)
|
||||||
|
return len(p), nil
|
||||||
|
}
|
||||||
Binary file not shown.
|
Before Width: | Height: | Size: 134 KiB After Width: | Height: | Size: 1.3 MiB |
Binary file not shown.
|
Before Width: | Height: | Size: 4.4 KiB After Width: | Height: | Size: 2.6 KiB |
Reference in New Issue
Block a user