Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
168 changes: 22 additions & 146 deletions cmd/eraser/cmd_monitor.go
Original file line number Diff line number Diff line change
Expand Up @@ -122,8 +122,7 @@ func runMonitor(days int, once bool, watch bool) error {
}

// scanInbox classifies and stores one inbox's broker replies, and with --watch
// keeps watching until ctx ends. Replies are attributed per broker via
// ResolveProfileForBroker, since a shared inbox serves several profiles.
// keeps watching until ctx ends.
func scanInbox(ctx context.Context, inboxCfg config.InboxConfig, brokerDB *broker.BrokerDatabase, store *history.Store, days int, once bool, watch bool) error {
monitor := inbox.NewMonitor(inboxCfg, brokerDB.Brokers)

Expand All @@ -132,109 +131,21 @@ func scanInbox(ctx context.Context, inboxCfg config.InboxConfig, brokerDB *broke
}
defer func() { _ = monitor.Disconnect() }()

fmt.Printf("📬 Monitoring %s for broker responses (last %d days)...\n", inboxCfg.Email, days)
fmt.Println()
fmt.Printf("📬 Monitoring %s for broker responses (last %d days)...\n\n", inboxCfg.Email, days)

emails, err := monitor.FetchBrokerEmails(ctx, days)
res, err := monitor.ScanAndStore(ctx, store, inbox.ScanOptions{Days: days})
if err != nil {
return fmt.Errorf("failed to fetch emails from %s: %w", inboxCfg.Email, err)
return err
}

if len(emails) == 0 {
fmt.Printf("No emails from known brokers found in %s.\n", inboxCfg.Email)
if !watch {
return nil
}
fmt.Printf("Found %d emails from data brokers in %s, %d new\n\n", res.Summary.Total, inboxCfg.Email, len(res.New))
for _, r := range res.New {
printClassifiedResponse(r)
}

// Classify and process each email
fmt.Printf("Found %d emails from data brokers in %s\n", len(emails), inboxCfg.Email)
fmt.Println()

var responses []inbox.ClassifiedResponse
for _, email := range emails {
classified := inbox.ClassifyResponse(&email)
responses = append(responses, classified)

profileID, err := store.ResolveProfileForBroker(email.BrokerID)
if err != nil {
profileID = history.DefaultProfileID
}

brokerResp := &history.BrokerResponse{
ProfileID: profileID,
BrokerID: email.BrokerID,
BrokerName: email.BrokerName,
ResponseType: string(classified.Type),
EmailFrom: email.From,
EmailSubject: email.Subject,
EmailBody: emailBody(email),
FormURL: classified.FormURL,
ConfirmURL: classified.ConfirmURL,
Confidence: classified.Confidence,
NeedsReview: classified.NeedsReview,
ReceivedAt: email.ReceivedAt,
}

inserted, err := store.AddBrokerResponseIfNew(brokerResp)
if err != nil {
fmt.Printf("⚠️ Failed to store response: %v\n", err)
} else if !inserted {
// Seen on an earlier scan. Re-applying its pipeline status would
// undo progress made since (e.g. a form you've since filled).
continue
}

// Update pipeline status for the broker
var pipelineStatus history.PipelineStatus
switch classified.Type {
case inbox.ResponseSuccess:
pipelineStatus = history.PipelineConfirmed
case inbox.ResponseFormRequired:
pipelineStatus = history.PipelineFormRequired
case inbox.ResponseConfirmationRequired:
pipelineStatus = history.PipelineAwaitingConfirmation
case inbox.ResponseRejected:
pipelineStatus = history.PipelineRejected
case inbox.ResponsePending:
pipelineStatus = history.PipelineAwaitingResponse
default:
pipelineStatus = history.PipelineAwaitingResponse
}

// Ignore error if no matching record
_ = store.UpdatePipelineStatus(profileID, email.BrokerID, pipelineStatus)

printClassifiedResponse(classified)
if res.Archived > 0 {
fmt.Printf("📁 Archived %d emails to '%s'\n", res.Archived, inboxCfg.ArchiveFolder)
}

// Archive processed emails if enabled
if inboxCfg.AutoArchive && len(emails) > 0 {
archiveFolder := inboxCfg.ArchiveFolder

// Ensure archive folder exists
if err := monitor.EnsureFolderExists(archiveFolder); err != nil {
fmt.Printf("⚠️ Could not create archive folder: %v\n", err)
} else {
// Collect UIDs to archive
var uidsToArchive []uint32
for _, email := range emails {
if email.UID > 0 {
uidsToArchive = append(uidsToArchive, email.UID)
}
}

if len(uidsToArchive) > 0 {
if err := monitor.ArchiveEmails(uidsToArchive, archiveFolder); err != nil {
fmt.Printf("⚠️ Could not archive emails: %v\n", err)
} else {
fmt.Printf("📁 Archived %d emails to '%s'\n", len(uidsToArchive), archiveFolder)
}
}
}
}

summary := inbox.SummarizeResponses(responses)
summary := res.Summary
fmt.Println()
fmt.Println("━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━")
fmt.Printf("📊 Summary for %s:\n", inboxCfg.Email)
Expand All @@ -247,50 +158,24 @@ func scanInbox(ctx context.Context, inboxCfg config.InboxConfig, brokerDB *broke
fmt.Printf(" ❓ Unknown: %d\n", summary.Unknown)
fmt.Printf(" 👁️ Need review: %d\n", summary.NeedReview)

if once {
if once || !watch {
return nil
}

if watch {
fmt.Println()
fmt.Printf("👀 Watching %s for new emails... (Ctrl+C to stop)\n", inboxCfg.Email)
err = monitor.WatchForNewEmails(ctx, func(email inbox.Email) {
fmt.Println()
fmt.Printf("👀 Watching %s for new emails... (Ctrl+C to stop)\n", inboxCfg.Email)

err := monitor.WatchForNewEmails(ctx, func(email inbox.Email) {
fmt.Println()
fmt.Printf("📨 New email from %s (%s)\n", email.BrokerName, email.From)

classified := inbox.ClassifyResponse(&email)
printClassifiedResponse(classified)

profileID, err := store.ResolveProfileForBroker(email.BrokerID)
if err != nil {
profileID = history.DefaultProfileID
}

brokerResp := &history.BrokerResponse{
ProfileID: profileID,
BrokerID: email.BrokerID,
BrokerName: email.BrokerName,
ResponseType: string(classified.Type),
EmailFrom: email.From,
EmailSubject: email.Subject,
EmailBody: emailBody(email),
FormURL: classified.FormURL,
ConfirmURL: classified.ConfirmURL,
Confidence: classified.Confidence,
NeedsReview: classified.NeedsReview,
ReceivedAt: email.ReceivedAt,
}
if _, err := store.AddBrokerResponseIfNew(brokerResp); err != nil {
fmt.Printf("⚠️ Failed to store response: %v\n", err)
}
})

if err != nil && err != context.Canceled {
return fmt.Errorf("watch error on %s: %w", inboxCfg.Email, err)
fmt.Printf("📨 New email from %s (%s)\n", email.BrokerName, email.From)
classified, _, err := inbox.RecordReply(store, &email, false)
if err != nil {
fmt.Printf("⚠️ Failed to store response: %v\n", err)
}
printClassifiedResponse(classified)
})
if err != nil && err != context.Canceled {
return fmt.Errorf("watch error on %s: %w", inboxCfg.Email, err)
}

return nil
}

Expand Down Expand Up @@ -324,12 +209,3 @@ func printClassifiedResponse(r inbox.ClassifiedResponse) {
fmt.Printf(" ⚠️ Confidence: %.0f%% - manual review recommended\n", r.Confidence*100)
}
}

// emailBody is the stored body: plain text, else HTML. Stored so later
// reclassification doesn't need to fetch the mail again.
func emailBody(e inbox.Email) string {
if e.Body != "" {
return e.Body
}
return e.HTMLBody
}
88 changes: 34 additions & 54 deletions cmd/eraser/cmd_send.go
Original file line number Diff line number Diff line change
Expand Up @@ -236,69 +236,49 @@ func runSend() error {
continue
}

emailMsg, err := tmplEngine.Render(cfg.Options.Template, activeProfile.Profile, b)
if cfg.Options.DryRun {
emailMsg, err := tmplEngine.Render(cfg.Options.Template, activeProfile.Profile, b)
if err != nil {
fmt.Printf(" ❌ Failed to render template: %v\n", err)
failCount++
continue
}
fmt.Printf(" 📧 Would send: %s\n", emailMsg.Subject)
fmt.Printf(" 📍 To: %s\n", b.Email)
successCount++
continue
}

record, err := email.SendRemoval(context.Background(), sender, tmplEngine, cfg.Options.Template, activeProfile, emailCfg.From, b)
if err != nil {
fmt.Printf(" ❌ Failed to render template: %v\n", err)
fmt.Printf(" ❌ %v\n", err)
failCount++
continue
}

if cfg.Options.DryRun {
fmt.Printf(" 📧 Would send: %s\n", emailMsg.Subject)
fmt.Printf(" 📍 To: %s\n", b.Email)
if record.Status == history.StatusSent {
fmt.Printf(" ✅ Sent successfully\n")
successCount++
authFails = 0
} else {
msg := email.Message{
To: b.Email,
From: emailCfg.From,
Subject: emailMsg.Subject,
Body: emailMsg.Body,
}

ctx := context.WithValue(context.Background(), email.SequenceKey, i)
result := sender.Send(ctx, msg)

// Record in history
record := &history.Record{
ProfileID: activeProfile.ID,
BrokerID: b.ID,
BrokerName: b.Name,
Email: b.Email,
Template: cfg.Options.Template,
SentAt: time.Now(),
}

if result.Success {
record.Status = history.StatusSent
record.MessageID = result.MessageID
fmt.Printf(" ✅ Sent successfully\n")
successCount++
authFails = 0
} else {
record.Status = history.StatusFailed
record.Error = result.Error.Error()
fmt.Printf(" ❌ Failed: %v\n", result.Error)
failCount++
if strings.Contains(strings.ToLower(record.Error), "auth") {
authFails++
}
}

if err := store.Add(record); err != nil {
fmt.Printf(" ⚠️ Failed to record history: %v\n", err)
fmt.Printf(" ❌ Failed: %s\n", record.Error)
failCount++
if strings.Contains(strings.ToLower(record.Error), "auth") {
authFails++
}
}
if err := store.Add(record); err != nil {
fmt.Printf(" ⚠️ Failed to record history: %v\n", err)
}

// Same cutoff as the web job sender: a bad password or a provider
// block would otherwise mark every due broker failed, and failed
// brokers are retried on every run.
if authFails >= 3 {
return fmt.Errorf("stopped after %d consecutive authentication failures (%d sent, %d failed) - check your email settings", authFails, successCount, failCount)
}
// Same cutoff as the web job sender: a bad password or a provider
// block would otherwise mark every due broker failed, and failed
// brokers are retried on every run.
if authFails >= 3 {
return fmt.Errorf("stopped after %d consecutive authentication failures (%d sent, %d failed) - check your email settings", authFails, successCount, failCount)
}

// Rate limiting
if i < len(brokers)-1 {
time.Sleep(time.Duration(cfg.Options.RateLimitMs) * time.Millisecond)
}
if i < len(brokers)-1 {
time.Sleep(time.Duration(cfg.Options.RateLimitMs) * time.Millisecond)
}
}

Expand Down
5 changes: 3 additions & 2 deletions docs/architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -28,10 +28,11 @@ eraser/
│ │ # validation (domain.go)
│ ├── config/config.go # User configuration (profile(s), email, options, inbox, pipeline)
│ ├── email/
│ │ ├── sender.go # Sender interface + NewSender (SMTP only)
│ │ ├── sender.go # NewSender (SMTP only), address validation
│ │ ├── removal.go # SendRemoval: render + send + history record, shared by CLI and web
│ │ └── smtp.go # SMTP implementation
│ ├── history/history.go # SQLite history tracking, pipeline status, per-profile scoping
│ ├── inbox/ # IMAP monitoring + reply classification (success/form-required/
│ ├── inbox/ # IMAP scan (scan.go: ScanAndStore/RecordReply) + reply classification (success/form-required/
│ │ # confirmation/rejection/pending/bounced)
│ ├── schedule/ # unattended cycles: shared lock + state file, launchd/systemd install
│ ├── template/
Expand Down
2 changes: 1 addition & 1 deletion docs/code-patterns.md
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@

Worth knowing before touching related code:

- `internal/web/handlers_jobs.go`'s `processSendJob` (the web UI's background sender) and `cmd/eraser/cmd_send.go`'s `send` command both need to respect `config.Options.DailySendLimit` - they used to disagree (the web UI had its own hardcoded, provider-based limit left over from the SendGrid/Resend era). Fixed, but if you add a third send path, wire it the same way.
- `internal/web/handlers_jobs.go`'s `processSendJob` (the web UI's background sender) and `cmd/eraser/cmd_send.go`'s `send` command both need to respect `config.Options.DailySendLimit` - they used to disagree (the web UI had its own hardcoded, provider-based limit left over from the SendGrid/Resend era). Fixed; both (and the single-broker web send) now go through `email.SendRemoval` for the render/send/record step, so a new send path should too and only add its own loop policy.
- `internal/web/server.go`'s `Server.config` is an `atomic.Pointer[config.Config]`, not a bare pointer - concurrent handlers and the background send-job goroutine read/write it. Always read via `s.getConfig()`; to write, load the current value, copy the struct (`newCfg := *cfg`), mutate the copy, then `s.config.Store(&newCfg)` - never mutate the struct returned by `getConfig()` in place, since another goroutine may be holding that exact pointer.
- The five by-ID methods on `internal/history/history.go`'s `Store` (`GetPendingTaskByID`, `CompletePendingTask`, `MarkTaskOpened`, `UpdateBrokerResponseClassification`, `UpdateBrokerResponseBody`) all take a `profileID string` and enforce it in the query (`AND profile_id = ?`). The getter returns `(nil, nil)` for a wrong-profile ID (same as genuinely not found); the four updaters are **silent no-ops** on a wrong-profile ID - they return `nil` either way, since a well-formed `UPDATE` matching zero rows isn't an error. If you need to know whether an update actually landed, re-fetch and check, don't trust the return value.
- `internal/browser/domain.go`'s `matchesAllowedDomain(rawURL string, allowedDomains []string) (bool, string, error)` is the one place domain-allowlist logic lives - `Browser.NavigateAndFill` and `ConfirmationHandler.ValidateDomain`/its redirect-hop check in `ClickConfirmationLink` all call it. If you add a new code path that navigates to or fetches an externally-supplied URL (e.g. from a parsed email), route it through this same allowlist rather than adding a one-off check - an empty/nil allowlist is treated as "no allowlist configured" (check passes), not "reject everything".
Expand Down
3 changes: 1 addition & 2 deletions docs/multi-profile.md
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,7 @@ Nearly every `Store` method takes a `profileID string` and filters by it. Two ar

### Shared inbox

A shared mailbox carries replies for every profile's sent requests together. When processing an inbound reply, attribute it to whichever profile actually emailed that broker - not to whatever profile happens to be "active" in the CLI/web session doing the scan. That's what `Store.ResolveProfileForBroker(brokerID string) (string, error)` is for: it looks up the most recent `removal_requests` row for that broker across *all* profiles and returns its `profile_id` (falling back to `"default"` if the broker was never emailed by anyone). Both `handleAPIInboxScan`/`handleAPIInboxRescan` (web) and `runMonitor` (CLI) call this per email before storing the classified response.
A shared mailbox carries replies for every profile's sent requests together. When processing an inbound reply, attribute it to whichever profile actually emailed that broker - not to whatever profile happens to be "active" in the CLI/web session doing the scan. That's what `Store.ResolveProfileForBroker(brokerID string) (string, error)` is for: it looks up the most recent `removal_requests` row for that broker across *all* profiles and returns its `profile_id` (falling back to `"default"` if the broker was never emailed by anyone). `inbox.RecordReply` (used by `Monitor.ScanAndStore` for the CLI monitor, the web scan/rescan and automated cycles) calls this per email before storing the classified response.

## Per-profile email accounts (`internal/config/config.go`)

Expand Down Expand Up @@ -94,7 +94,6 @@ just `inbox`.
- `renderWithCSRF` injects `Profiles`/`ActiveProfile`/`CurrentPath` into every page's template data, so `layout.html`'s nav can render the switcher unconditionally without every handler wiring it manually
- The switcher itself is a `<select>` inside a small auto-submitting `<form>` (desktop nav + mobile nav), only rendered when `len(.Profiles) > 1`
- Background send jobs (`internal/web/job.go`) carry `ProfileID`; `JobManager.GetActive(profileID)` and `.Create(total, profileID)` are profile-scoped, so two profiles can have a send running concurrently without colliding on "job already active"
- `PersistentJobState` (the on-restart job-resume format) also carries `ProfileID`, so a resumed job after a server restart is attributed correctly

## Adding a new per-profile data path

Expand Down
6 changes: 2 additions & 4 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -4,20 +4,18 @@ go 1.26

require (
filippo.io/csrf v0.2.1
github.com/PuerkitoBio/goquery v1.13.0
github.com/chromedp/chromedp v0.16.0
github.com/emersion/go-imap v1.2.1
github.com/emersion/go-message v0.18.2
github.com/go-chi/chi/v5 v5.3.2
github.com/google/uuid v1.6.0
github.com/spf13/cobra v1.10.2
golang.org/x/net v0.58.0
golang.org/x/sys v0.47.0
gopkg.in/yaml.v3 v3.0.1
modernc.org/sqlite v1.59.0
)

require (
github.com/andybalholm/cascadia v1.3.4 // indirect
github.com/chromedp/cdproto v0.0.0-20260804232424-e85f50dbfd32 // indirect
github.com/chromedp/sysutil v1.1.0 // indirect
github.com/dustin/go-humanize v1.0.1 // indirect
Expand All @@ -26,13 +24,13 @@ require (
github.com/gobwas/httphead v0.1.0 // indirect
github.com/gobwas/pool v0.2.1 // indirect
github.com/gobwas/ws v1.4.0 // indirect
github.com/google/uuid v1.6.0 // indirect
github.com/inconshreveable/mousetrap v1.1.0 // indirect
github.com/mattn/go-isatty v0.0.24 // indirect
github.com/ncruces/go-strftime v1.0.0 // indirect
github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec // indirect
github.com/spf13/pflag v1.0.10 // indirect
golang.org/x/mod v0.40.0 // indirect
golang.org/x/net v0.58.0 // indirect
golang.org/x/text v0.41.0 // indirect
golang.org/x/tools v0.49.0 // indirect
modernc.org/libc v1.75.7 // indirect
Expand Down
4 changes: 0 additions & 4 deletions go.sum
Original file line number Diff line number Diff line change
@@ -1,9 +1,5 @@
filippo.io/csrf v0.2.1 h1:MdV/y9xOECwJko48lPkH9NaYNpZ6kYfaNlgnZk4k6Uo=
filippo.io/csrf v0.2.1/go.mod h1:eVfdeENlqr/ErpNx4E5I6a11I1aP0WL/PPkzKD1d960=
github.com/PuerkitoBio/goquery v1.13.0 h1:mqHbjD7Jmnul4DTR24LKTjo1uUmHUh072kteGV+xpFM=
github.com/PuerkitoBio/goquery v1.13.0/go.mod h1:Hip5mdBL8K2wEGKJdr27sRaNwIdDajmCwB/ExUPwW+g=
github.com/andybalholm/cascadia v1.3.4 h1:vM2lgh0Vru9Vwyfm4cQqWP2HHMW0u0+2PAW7Q38Qufg=
github.com/andybalholm/cascadia v1.3.4/go.mod h1:BLRmbRjpEtNKieZOCCvYj4RqN+KRA41GBe/5O+G93kM=
github.com/chromedp/cdproto v0.0.0-20260804232424-e85f50dbfd32 h1:6JI+JS7Zef+bMzZQ+OgzTHf79v3GqdvP6rD0FaP9CMk=
github.com/chromedp/cdproto v0.0.0-20260804232424-e85f50dbfd32/go.mod h1:RwFsSODCtFExll+GhHM6R92SARHR3Z3oipaxLHj46C0=
github.com/chromedp/chromedp v0.16.0 h1:rOO4deOm4CbZgBCa8mD9g2rDyIoNs0BkgvNrlbp5ouk=
Expand Down
Loading
Loading