diff --git a/cmd/eraser/cmd_monitor.go b/cmd/eraser/cmd_monitor.go index 67e0aef..19cac20 100644 --- a/cmd/eraser/cmd_monitor.go +++ b/cmd/eraser/cmd_monitor.go @@ -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) @@ -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) @@ -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 } @@ -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 -} diff --git a/cmd/eraser/cmd_send.go b/cmd/eraser/cmd_send.go index 2700b47..b767f7e 100644 --- a/cmd/eraser/cmd_send.go +++ b/cmd/eraser/cmd_send.go @@ -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) } } diff --git a/docs/architecture.md b/docs/architecture.md index 85c5276..9b64197 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -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/ diff --git a/docs/code-patterns.md b/docs/code-patterns.md index 6b09762..302ea16 100644 --- a/docs/code-patterns.md +++ b/docs/code-patterns.md @@ -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". diff --git a/docs/multi-profile.md b/docs/multi-profile.md index 3b717ef..d6439f8 100644 --- a/docs/multi-profile.md +++ b/docs/multi-profile.md @@ -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`) @@ -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 ``, ""}, + {"recaptcha v2", `
`, CaptchaTypeRecaptchaV2}, + {"recaptcha v3", ``, CaptchaTypeRecaptchaV3}, + {"hcaptcha", `
`, CaptchaTypeHCaptcha}, + {"turnstile", `
`, CaptchaTypeTurnstile}, + {"funcaptcha", `
`, CaptchaTypeFunCaptcha}, + {"cloudflare", `Just a moment...

checking

`, CaptchaTypeCloudflare}, + {"text captcha", `

Enter the code shown

`, CaptchaTypeTextCaptcha}, + {"image captcha", `x`, CaptchaTypeImageCaptcha}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + ctx, cancel := context.WithTimeout(b.ctx, b.config.Timeout) + defer cancel() + if err := chromedp.Run(ctx, chromedp.Navigate("data:text/html,"+tc.body+"")); err != nil { + t.Fatalf("navigate: %v", err) + } + got, err := b.detectCaptcha(ctx) + if err != nil { + t.Fatalf("detectCaptcha: %v", err) + } + if got.Type != tc.want || got.Found != (tc.want != "") { + t.Fatalf("got %+v, want type %q", got, tc.want) + } + }) + } +} diff --git a/internal/browser/captcha.go b/internal/browser/captcha.go index 72535bb..bf8b59a 100644 --- a/internal/browser/captcha.go +++ b/internal/browser/captcha.go @@ -6,17 +6,12 @@ import ( "github.com/chromedp/chromedp" ) -// CaptchaInfo contains information about detected CAPTCHAs +// CaptchaInfo is the result of detectCaptcha. type CaptchaInfo struct { - Found bool - Type string - FrameSrc string - ElementID string - Confidence float64 - Description string + Found bool + Type string // one of the CaptchaType constants } -// CaptchaType constants const ( CaptchaTypeRecaptchaV2 = "recaptcha_v2" CaptchaTypeRecaptchaV3 = "recaptcha_v3" @@ -29,413 +24,66 @@ const ( CaptchaTypeUnknown = "unknown" ) +// detectCaptchaJS returns the first CAPTCHA type found on the page, checked +// most specific first, or "" for none. +const detectCaptchaJS = `(function() { + var q = function(sel) { return document.querySelector(sel); }; + var body = document.body ? document.body.innerText.toLowerCase() : ''; + + if (q('iframe[src*="recaptcha"]') || q('.g-recaptcha, [data-sitekey]') || typeof grecaptcha !== 'undefined') return 'recaptcha_v2'; + var scripts = document.querySelectorAll('script[src*="recaptcha"]'); + for (var i = 0; i < scripts.length; i++) { + if (scripts[i].src.includes('render=')) return 'recaptcha_v3'; + } + if (q('script[src*="enterprise.js"]')) return 'recaptcha_v3'; + if (q('iframe[src*="hcaptcha"]') || q('.h-captcha, [data-hcaptcha-sitekey]') || typeof hcaptcha !== 'undefined') return 'hcaptcha'; + if (q('iframe[src*="challenges.cloudflare.com"]') || q('.cf-turnstile, [data-turnstile-sitekey]')) return 'cloudflare_turnstile'; + if (q('iframe[src*="funcaptcha"], iframe[src*="arkoselabs"]')) return 'funcaptcha'; + var fc = q('#FunCaptcha, [data-callback]'); + if (fc && fc.id && fc.id.toLowerCase().includes('funcaptcha')) return 'funcaptcha'; + + var title = document.title.toLowerCase(); + if (title.includes('just a moment') || title.includes('checking your browser') || q('form#challenge-form') || + q('.ray-id, [data-ray]') || (body.includes('ray id') && body.includes('cloudflare'))) return 'cloudflare_challenge'; + + var keywords = ['captcha', 'verification code', 'security code', 'enter the code', 'type the characters', + 'verify you are human', 'prove you are human', 'i am not a robot', 'human verification']; + var img = q('img[src*="captcha"], img[alt*="captcha"], .captcha-image'); + var input = q('input[name*="captcha"], input[id*="captcha"], input[placeholder*="captcha" i]'); + if ((img || input) && keywords.some(function(k) { return body.includes(k); })) return input ? 'text_captcha' : 'unknown'; + if (q('img[src*="captcha"], img[alt*="captcha"]')) return 'image_captcha'; + return ''; +})()` + // detectCaptcha errors only on context failure (browser died or timed out), // never for "no CAPTCHA". func (b *Browser) detectCaptcha(ctx context.Context) (CaptchaInfo, error) { - // Check for reCAPTCHA v2 - if result, err := detectRecaptchaV2(ctx); err != nil { - return CaptchaInfo{}, err - } else if result.Found { - return result, nil - } - - // Check for reCAPTCHA v3 (invisible) - if result, err := detectRecaptchaV3(ctx); err != nil { - return CaptchaInfo{}, err - } else if result.Found { - return result, nil - } - - // Check for hCaptcha - if result, err := detectHCaptcha(ctx); err != nil { - return CaptchaInfo{}, err - } else if result.Found { - return result, nil - } - - // Check for Cloudflare Turnstile - if result, err := detectTurnstile(ctx); err != nil { - return CaptchaInfo{}, err - } else if result.Found { - return result, nil - } - - // Check for FunCaptcha - if result, err := detectFunCaptcha(ctx); err != nil { - return CaptchaInfo{}, err - } else if result.Found { - return result, nil - } - - // Check for Cloudflare challenge page - if result, err := detectCloudflareChallenge(ctx); err != nil { - return CaptchaInfo{}, err - } else if result.Found { - return result, nil - } - - // Check for generic image/text CAPTCHA - if result, err := detectGenericCaptcha(ctx); err != nil { - return CaptchaInfo{}, err - } else if result.Found { - return result, nil - } - - return CaptchaInfo{}, nil -} - -// detectRecaptchaV2 checks for Google reCAPTCHA v2 -func detectRecaptchaV2(ctx context.Context) (CaptchaInfo, error) { - js := `(function() { - // Check for reCAPTCHA iframe - var iframe = document.querySelector('iframe[src*="recaptcha"]'); - if (iframe) { - return {found: true, src: iframe.src, type: 'iframe'}; - } - - // Check for reCAPTCHA div - var div = document.querySelector('.g-recaptcha, [data-sitekey]'); - if (div) { - return {found: true, id: div.id || '', type: 'div'}; - } - - // Check for grecaptcha object - if (typeof grecaptcha !== 'undefined') { - return {found: true, type: 'script'}; - } - - return {found: false}; - })()` - - var result map[string]interface{} - err := chromedp.Run(ctx, chromedp.Evaluate(js, &result)) - if err != nil { - if isContextErr(err) { - return CaptchaInfo{}, err - } - return CaptchaInfo{}, nil - } - - found, _ := result["found"].(bool) - if !found { - return CaptchaInfo{}, nil - } - - info := CaptchaInfo{ - Found: true, - Type: CaptchaTypeRecaptchaV2, - Confidence: 0.95, - Description: "Google reCAPTCHA v2 detected", - } - - if src, ok := result["src"].(string); ok { - info.FrameSrc = src - } - if id, ok := result["id"].(string); ok { - info.ElementID = id - } - - return info, nil -} - -// detectRecaptchaV3 checks for Google reCAPTCHA v3 (invisible) -func detectRecaptchaV3(ctx context.Context) (CaptchaInfo, error) { - js := `(function() { - // reCAPTCHA v3 is typically loaded via script - var scripts = document.querySelectorAll('script[src*="recaptcha"]'); - for (var i = 0; i < scripts.length; i++) { - if (scripts[i].src.includes('render=')) { - return {found: true, src: scripts[i].src}; - } - } - - // Check for enterprise version - scripts = document.querySelectorAll('script[src*="enterprise.js"]'); - if (scripts.length > 0) { - return {found: true, type: 'enterprise'}; - } - - return {found: false}; - })()` - - var result map[string]interface{} - err := chromedp.Run(ctx, chromedp.Evaluate(js, &result)) - if err != nil { - if isContextErr(err) { - return CaptchaInfo{}, err - } - return CaptchaInfo{}, nil - } - - found, _ := result["found"].(bool) - if !found { - return CaptchaInfo{}, nil - } - - return CaptchaInfo{ - Found: true, - Type: CaptchaTypeRecaptchaV3, - Confidence: 0.85, - Description: "Google reCAPTCHA v3 (invisible) detected - may not require interaction", - }, nil -} - -// detectHCaptcha checks for hCaptcha -func detectHCaptcha(ctx context.Context) (CaptchaInfo, error) { - js := `(function() { - // Check for hCaptcha iframe - var iframe = document.querySelector('iframe[src*="hcaptcha"]'); - if (iframe) { - return {found: true, src: iframe.src}; - } - - // Check for hCaptcha div - var div = document.querySelector('.h-captcha, [data-hcaptcha-sitekey]'); - if (div) { - return {found: true, id: div.id || ''}; - } - - // Check for hcaptcha object - if (typeof hcaptcha !== 'undefined') { - return {found: true, type: 'script'}; - } - - return {found: false}; - })()` - - var result map[string]interface{} - err := chromedp.Run(ctx, chromedp.Evaluate(js, &result)) - if err != nil { - if isContextErr(err) { - return CaptchaInfo{}, err - } - return CaptchaInfo{}, nil - } - - found, _ := result["found"].(bool) - if !found { - return CaptchaInfo{}, nil - } - - info := CaptchaInfo{ - Found: true, - Type: CaptchaTypeHCaptcha, - Confidence: 0.95, - Description: "hCaptcha detected", - } - - if src, ok := result["src"].(string); ok { - info.FrameSrc = src - } - - return info, nil -} - -// detectTurnstile checks for Cloudflare Turnstile -func detectTurnstile(ctx context.Context) (CaptchaInfo, error) { - js := `(function() { - // Check for Turnstile iframe - var iframe = document.querySelector('iframe[src*="challenges.cloudflare.com"]'); - if (iframe) { - return {found: true, src: iframe.src}; - } - - // Check for Turnstile div - var div = document.querySelector('.cf-turnstile, [data-turnstile-sitekey]'); - if (div) { - return {found: true, id: div.id || ''}; - } - - return {found: false}; - })()` - - var result map[string]interface{} - err := chromedp.Run(ctx, chromedp.Evaluate(js, &result)) - if err != nil { - if isContextErr(err) { - return CaptchaInfo{}, err - } - return CaptchaInfo{}, nil - } - - found, _ := result["found"].(bool) - if !found { - return CaptchaInfo{}, nil - } - - return CaptchaInfo{ - Found: true, - Type: CaptchaTypeTurnstile, - Confidence: 0.95, - Description: "Cloudflare Turnstile detected", - }, nil -} - -// detectFunCaptcha checks for Arkose Labs FunCaptcha -func detectFunCaptcha(ctx context.Context) (CaptchaInfo, error) { - js := `(function() { - // Check for FunCaptcha iframe - var iframe = document.querySelector('iframe[src*="funcaptcha"], iframe[src*="arkoselabs"]'); - if (iframe) { - return {found: true, src: iframe.src}; - } - - // Check for FunCaptcha div - var div = document.querySelector('#FunCaptcha, [data-callback]'); - if (div && div.id && div.id.toLowerCase().includes('funcaptcha')) { - return {found: true, id: div.id}; - } - - return {found: false}; - })()` - - var result map[string]interface{} - err := chromedp.Run(ctx, chromedp.Evaluate(js, &result)) - if err != nil { - if isContextErr(err) { - return CaptchaInfo{}, err - } - return CaptchaInfo{}, nil - } - - found, _ := result["found"].(bool) - if !found { - return CaptchaInfo{}, nil - } - - return CaptchaInfo{ - Found: true, - Type: CaptchaTypeFunCaptcha, - Confidence: 0.90, - Description: "Arkose Labs FunCaptcha detected", - }, nil -} - -// detectCloudflareChallenge checks for Cloudflare challenge page -func detectCloudflareChallenge(ctx context.Context) (CaptchaInfo, error) { - js := `(function() { - // Check for Cloudflare challenge indicators - var title = document.title.toLowerCase(); - if (title.includes('just a moment') || title.includes('checking your browser')) { - return {found: true, type: 'title'}; - } - - // Check for challenge form - var form = document.querySelector('form#challenge-form'); - if (form) { - return {found: true, type: 'form'}; - } - - // Check for ray ID (Cloudflare identifier) - var rayId = document.querySelector('.ray-id, [data-ray]'); - var body = document.body.innerText.toLowerCase(); - if (rayId || (body.includes('ray id') && body.includes('cloudflare'))) { - return {found: true, type: 'rayid'}; - } - - return {found: false}; - })()` - - var result map[string]interface{} - err := chromedp.Run(ctx, chromedp.Evaluate(js, &result)) - if err != nil { + var captchaType string + if err := chromedp.Run(ctx, chromedp.Evaluate(detectCaptchaJS, &captchaType)); err != nil { if isContextErr(err) { return CaptchaInfo{}, err } - return CaptchaInfo{}, nil + return CaptchaInfo{}, nil // page script failed: treat as no CAPTCHA } - - found, _ := result["found"].(bool) - if !found { - return CaptchaInfo{}, nil - } - - return CaptchaInfo{ - Found: true, - Type: CaptchaTypeCloudflare, - Confidence: 0.90, - Description: "Cloudflare challenge page detected - wait or solve challenge", - }, nil + return CaptchaInfo{Found: captchaType != "", Type: captchaType}, nil } -// detectGenericCaptcha checks for generic image/text CAPTCHAs -func detectGenericCaptcha(ctx context.Context) (CaptchaInfo, error) { - js := `(function() { - var body = document.body.innerHTML.toLowerCase(); - var text = document.body.innerText.toLowerCase(); - - // Check for CAPTCHA-related keywords - var keywords = ['captcha', 'verification code', 'security code', 'enter the code', - 'type the characters', 'verify you are human', 'prove you are human', - 'i am not a robot', 'human verification']; - - for (var i = 0; i < keywords.length; i++) { - if (text.includes(keywords[i])) { - // Look for associated input and image - var img = document.querySelector('img[src*="captcha"], img[alt*="captcha"], .captcha-image'); - var input = document.querySelector('input[name*="captcha"], input[id*="captcha"], input[placeholder*="captcha" i]'); - - if (img || input) { - return {found: true, keyword: keywords[i], hasImage: !!img, hasInput: !!input}; - } - } - } - - // Check for CAPTCHA images even without keywords - var captchaImg = document.querySelector('img[src*="captcha"], img[alt*="captcha"]'); - if (captchaImg) { - return {found: true, type: 'image', src: captchaImg.src}; - } - - return {found: false}; - })()` - - var result map[string]interface{} - err := chromedp.Run(ctx, chromedp.Evaluate(js, &result)) - if err != nil { - if isContextErr(err) { - return CaptchaInfo{}, err - } - return CaptchaInfo{}, nil - } - - found, _ := result["found"].(bool) - if !found { - return CaptchaInfo{}, nil - } - - captchaType := CaptchaTypeUnknown - if imgType, ok := result["type"].(string); ok && imgType == "image" { - captchaType = CaptchaTypeImageCaptcha - } else if hasInput, ok := result["hasInput"].(bool); ok && hasInput { - captchaType = CaptchaTypeTextCaptcha - } - - description := "Generic CAPTCHA detected" - if keyword, ok := result["keyword"].(string); ok { - description = "CAPTCHA detected: " + keyword - } - - return CaptchaInfo{ - Found: true, - Type: captchaType, - Confidence: 0.75, - Description: description, - }, nil -} - -// IsCaptchaBlocking returns true if the CAPTCHA requires human intervention +// IsCaptchaBlocking returns true if the CAPTCHA requires human intervention. +// reCAPTCHA v3 is invisible and usually doesn't. func (c CaptchaInfo) IsCaptchaBlocking() bool { - if !c.Found { - return false - } - - // reCAPTCHA v3 is invisible and may not block - if c.Type == CaptchaTypeRecaptchaV3 { - return false - } + return c.Found && c.Type != CaptchaTypeRecaptchaV3 +} - return true +var captchaDescriptions = map[string]string{ + CaptchaTypeRecaptchaV2: "Google reCAPTCHA v2 - Click the checkbox and/or solve image puzzles", + CaptchaTypeRecaptchaV3: "Google reCAPTCHA v3 - Usually invisible, may auto-pass", + CaptchaTypeHCaptcha: "hCaptcha - Select images matching the description", + CaptchaTypeTurnstile: "Cloudflare Turnstile - Usually auto-passes after brief check", + CaptchaTypeFunCaptcha: "FunCaptcha - Complete interactive puzzles", + CaptchaTypeImageCaptcha: "Image CAPTCHA - Type the characters shown in the image", + CaptchaTypeTextCaptcha: "Text CAPTCHA - Enter the verification code", + CaptchaTypeCloudflare: "Cloudflare Challenge - Wait or complete verification", + CaptchaTypeUnknown: "Unknown CAPTCHA type - Manual inspection required", } // GetCaptchaDescription returns a human-readable description @@ -443,22 +91,8 @@ func (c CaptchaInfo) GetCaptchaDescription() string { if !c.Found { return "No CAPTCHA detected" } - - descriptions := map[string]string{ - CaptchaTypeRecaptchaV2: "Google reCAPTCHA v2 - Click the checkbox and/or solve image puzzles", - CaptchaTypeRecaptchaV3: "Google reCAPTCHA v3 - Usually invisible, may auto-pass", - CaptchaTypeHCaptcha: "hCaptcha - Select images matching the description", - CaptchaTypeTurnstile: "Cloudflare Turnstile - Usually auto-passes after brief check", - CaptchaTypeFunCaptcha: "FunCaptcha - Complete interactive puzzles", - CaptchaTypeImageCaptcha: "Image CAPTCHA - Type the characters shown in the image", - CaptchaTypeTextCaptcha: "Text CAPTCHA - Enter the verification code", - CaptchaTypeCloudflare: "Cloudflare Challenge - Wait or complete verification", - CaptchaTypeUnknown: "Unknown CAPTCHA type - Manual inspection required", - } - - if desc, ok := descriptions[c.Type]; ok { + if desc, ok := captchaDescriptions[c.Type]; ok { return desc } - - return c.Description + return captchaDescriptions[CaptchaTypeUnknown] } diff --git a/internal/browser/captcha_test.go b/internal/browser/captcha_test.go index 477a862..0a45027 100644 --- a/internal/browser/captcha_test.go +++ b/internal/browser/captcha_test.go @@ -3,7 +3,7 @@ package browser import "testing" // IsCaptchaBlocking and GetCaptchaDescription are pure methods on an -// already-populated CaptchaInfo - unlike the detectXxx family (which drive a +// already-populated CaptchaInfo - unlike detectCaptcha (which drives a // live browser via chromedp and need a real Chrome instance, see // browser_chrome_test.go), these are cheap to check directly. @@ -45,11 +45,7 @@ func TestCaptchaInfoGetCaptchaDescription(t *testing.T) { {"not found", CaptchaInfo{Found: false}, "No CAPTCHA detected"}, {"recaptcha v2 has a known description", CaptchaInfo{Found: true, Type: CaptchaTypeRecaptchaV2}, "Google reCAPTCHA v2 - Click the checkbox and/or solve image puzzles"}, {"turnstile has a known description", CaptchaInfo{Found: true, Type: CaptchaTypeTurnstile}, "Cloudflare Turnstile - Usually auto-passes after brief check"}, - { - "unrecognized type falls back to the detector's own Description", - CaptchaInfo{Found: true, Type: "some_future_captcha_type", Description: "custom detector description"}, - "custom detector description", - }, + {"unrecognized type falls back to unknown", CaptchaInfo{Found: true, Type: "some_future_captcha_type"}, "Unknown CAPTCHA type - Manual inspection required"}, } for _, tt := range tests { diff --git a/internal/email/removal.go b/internal/email/removal.go new file mode 100644 index 0000000..b42eebe --- /dev/null +++ b/internal/email/removal.go @@ -0,0 +1,43 @@ +package email + +import ( + "context" + "fmt" + "time" + + "github.com/drumandbytes/eraser/internal/broker" + "github.com/drumandbytes/eraser/internal/config" + "github.com/drumandbytes/eraser/internal/history" + "github.com/drumandbytes/eraser/internal/template" +) + +// SendRemoval renders the removal request for b, sends it from `from`, and +// returns the history record for the attempt (sent or failed) for the caller +// to store. A render failure returns an error and no record: nothing was +// attempted. +func SendRemoval(ctx context.Context, s *SMTPSender, eng *template.Engine, tmpl string, np config.NamedProfile, from string, b broker.Broker) (*history.Record, error) { + rendered, err := eng.Render(tmpl, np.Profile, b) + if err != nil { + return nil, fmt.Errorf("failed to render template: %w", err) + } + result := s.Send(ctx, Message{To: b.Email, From: from, Subject: rendered.Subject, Body: rendered.Body}) + + record := &history.Record{ + ProfileID: np.ID, + BrokerID: b.ID, + BrokerName: b.Name, + Email: b.Email, + Template: tmpl, + SentAt: time.Now(), + Status: history.StatusSent, + MessageID: result.MessageID, + } + if !result.Success { + record.Status = history.StatusFailed + record.Error = "unknown error" + if result.Error != nil { + record.Error = result.Error.Error() + } + } + return record, nil +} diff --git a/internal/email/removal_test.go b/internal/email/removal_test.go new file mode 100644 index 0000000..0ace3d8 --- /dev/null +++ b/internal/email/removal_test.go @@ -0,0 +1,53 @@ +package email + +import ( + "context" + "net" + "strconv" + "strings" + "testing" + + "github.com/drumandbytes/eraser/internal/broker" + "github.com/drumandbytes/eraser/internal/config" + "github.com/drumandbytes/eraser/internal/history" + "github.com/drumandbytes/eraser/internal/template" +) + +func TestSendRemoval(t *testing.T) { + eng, err := template.NewEngine() + if err != nil { + t.Fatal(err) + } + np := config.NamedProfile{ID: "jane", Profile: config.Profile{FirstName: "Jane", LastName: "Doe", Email: "jane@example.org"}} + b := broker.Broker{ID: "acme", Name: "Acme", Email: "privacy@acme.example"} + + addr, data := recordingSMTPServer(t) + host, portStr, _ := net.SplitHostPort(addr) + port, _ := strconv.Atoi(portStr) + s := NewSMTPSender(config.SMTPConfig{Host: host, Port: port}, "jane@example.org") + + rec, err := SendRemoval(context.Background(), s, eng, "gdpr", np, "jane@example.org", b) + if err != nil { + t.Fatal(err) + } + if rec.Status != history.StatusSent || rec.ProfileID != "jane" || rec.BrokerID != "acme" || rec.Email != b.Email || rec.Template != "gdpr" || rec.MessageID == "" { + t.Fatalf("unexpected record: %+v", rec) + } + if msg := <-data; !strings.Contains(msg, "Jane") { + t.Fatalf("rendered request not sent:\n%s", msg) + } + + // Nothing listening: a failed attempt is still a record, with the error. + dead := NewSMTPSender(config.SMTPConfig{Host: "127.0.0.1", Port: 1}, "jane@example.org") + rec, err = SendRemoval(context.Background(), dead, eng, "gdpr", np, "jane@example.org", b) + if err != nil { + t.Fatal(err) + } + if rec.Status != history.StatusFailed || rec.Error == "" || rec.MessageID != "" { + t.Fatalf("failed send recorded as %+v", rec) + } + + if _, err := SendRemoval(context.Background(), s, eng, "no-such-template", np, "jane@example.org", b); err == nil { + t.Fatal("unknown template: want error, got nil") + } +} diff --git a/internal/email/sender.go b/internal/email/sender.go index 15d1206..d8313ae 100644 --- a/internal/email/sender.go +++ b/internal/email/sender.go @@ -8,16 +8,6 @@ import ( "github.com/drumandbytes/eraser/internal/config" ) -// ctxKey is an unexported type for context values defined by this package, -// so its keys can never collide with keys from another package using the -// same underlying string (see https://pkg.go.dev/context#WithValue). -type ctxKey string - -// SequenceKey is the context key runSend() uses to pass each broker's -// position in the current batch through to a Sender, so implementations -// (e.g. the SMTP sender) can fold it into a generated message ID. -const SequenceKey ctxKey = "sequence" - type Message struct { To string From string diff --git a/internal/email/smtp.go b/internal/email/smtp.go index 115fba7..c6c32e9 100644 --- a/internal/email/smtp.go +++ b/internal/email/smtp.go @@ -2,9 +2,11 @@ package email import ( "context" + "crypto/rand" "crypto/tls" "fmt" "net" + "net/mail" "net/smtp" "strings" @@ -31,7 +33,16 @@ func (s *SMTPSender) Send(ctx context.Context, msg Message) Result { addr := fmt.Sprintf("%s:%d", s.config.Host, s.config.Port) + // Our own Message-ID, so the one recorded in history is the one the + // broker actually receives (and quotes back in In-Reply-To). + domain := "localhost" + if from, err := mail.ParseAddress(msg.From); err == nil { + domain = from.Address[strings.LastIndex(from.Address, "@")+1:] + } + messageID := "<" + rand.Text() + "@" + domain + ">" + var message strings.Builder + fmt.Fprintf(&message, "Message-ID: %s\r\n", messageID) fmt.Fprintf(&message, "From: %s\r\n", msg.From) fmt.Fprintf(&message, "To: %s\r\n", msg.To) fmt.Fprintf(&message, "Subject: %s\r\n", msg.Subject) @@ -59,7 +70,7 @@ func (s *SMTPSender) Send(ctx context.Context, msg Message) Result { return Result{ Success: true, - MessageID: fmt.Sprintf("smtp-%s-%d", msg.To, ctx.Value(SequenceKey)), + MessageID: messageID, } } diff --git a/internal/email/smtp_test.go b/internal/email/smtp_test.go index f923b47..8ff5991 100644 --- a/internal/email/smtp_test.go +++ b/internal/email/smtp_test.go @@ -1,9 +1,11 @@ package email import ( + "bufio" "context" "net" "strconv" + "strings" "testing" "time" @@ -90,3 +92,71 @@ func splitHostPortForTest(t *testing.T, addr string) (string, int) { } return host, port } + +// recordingSMTPServer speaks just enough plaintext SMTP to accept one +// message and hands back its DATA. +func recordingSMTPServer(t *testing.T) (addr string, data <-chan string) { + t.Helper() + ln, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatalf("listen: %v", err) + } + t.Cleanup(func() { _ = ln.Close() }) + out := make(chan string, 1) + go func() { + conn, err := ln.Accept() + if err != nil { + return + } + defer func() { _ = conn.Close() }() + r := bufio.NewReader(conn) + reply := func(s string) { _, _ = conn.Write([]byte(s + "\r\n")) } + reply("220 test") + for { + line, err := r.ReadString('\n') + if err != nil { + return + } + switch cmd := strings.ToUpper(strings.TrimSpace(line)); { + case strings.HasPrefix(cmd, "EHLO"), strings.HasPrefix(cmd, "HELO"): + reply("250 test") + case strings.HasPrefix(cmd, "DATA"): + reply("354 go ahead") + var b strings.Builder + for { + l, err := r.ReadString('\n') + if err != nil || l == ".\r\n" { + break + } + b.WriteString(l) + } + out <- b.String() + reply("250 ok") + case strings.HasPrefix(cmd, "QUIT"): + reply("221 bye") + return + default: + reply("250 ok") + } + } + }() + return ln.Addr().String(), out +} + +func TestSendRecordsTheMessageIDItSends(t *testing.T) { + addr, data := recordingSMTPServer(t) + host, portStr, _ := net.SplitHostPort(addr) + port, _ := strconv.Atoi(portStr) + s := NewSMTPSender(config.SMTPConfig{Host: host, Port: port}, "Jane Doe ") + + res := s.Send(context.Background(), Message{To: "privacy@acme.example", From: "Jane Doe ", Subject: "Erasure request", Body: "hi"}) + if !res.Success { + t.Fatalf("send failed: %v", res.Error) + } + if !strings.HasPrefix(res.MessageID, "<") || !strings.HasSuffix(res.MessageID, "@example.org>") { + t.Fatalf("MessageID %q isn't ", res.MessageID) + } + if got := <-data; !strings.Contains(got, "Message-ID: "+res.MessageID+"\r\n") { + t.Fatalf("sent message lacks header Message-ID: %s:\n%s", res.MessageID, got) + } +} diff --git a/internal/history/history.go b/internal/history/history.go index cc3c0fd..25f4743 100644 --- a/internal/history/history.go +++ b/internal/history/history.go @@ -730,59 +730,76 @@ func (s *Store) AddBrokerResponseIfNew(resp *BrokerResponse) (inserted bool, err return true, s.AddBrokerResponse(resp) } -func (s *Store) GetBrokerResponseByID(id int64, profileID string) (*BrokerResponse, error) { - query := `SELECT id, profile_id, broker_id, broker_name, response_type, email_from, email_subject, email_body, - form_url, confirm_url, confidence, needs_review, received_at, processed_at, created_at - FROM broker_responses WHERE id = ? AND profile_id = ?` +// Column lists for scanBrokerResponse. The no-body variant is for listing +// pages, which don't show bodies and can hold hundreds of rows. +const ( + brokerResponseCols = `id, profile_id, broker_id, broker_name, response_type, email_from, email_subject, email_body, + form_url, confirm_url, confidence, needs_review, received_at, processed_at, created_at` + brokerResponseColsNoBody = `id, profile_id, broker_id, broker_name, response_type, email_from, email_subject, '' AS email_body, + form_url, confirm_url, confidence, needs_review, received_at, processed_at, created_at` +) +func scanBrokerResponse(scanner interface{ Scan(...any) error }) (*BrokerResponse, error) { var r BrokerResponse var needsReviewInt int - var receivedAtStr, processedAtStr, createdAtStr sql.NullString + var receivedAt, processedAt, createdAt sql.NullString var emailBody, formURL, confirmURL sql.NullString - if err := s.db.QueryRow(query, id, normalizeProfileID(profileID)).Scan( - &r.ID, &r.ProfileID, &r.BrokerID, &r.BrokerName, &r.ResponseType, &r.EmailFrom, &r.EmailSubject, &emailBody, - &formURL, &confirmURL, &r.Confidence, &needsReviewInt, &receivedAtStr, &processedAtStr, &createdAtStr); err != nil { - if err == sql.ErrNoRows { - return nil, nil - } - return nil, fmt.Errorf("failed to get broker response: %w", err) + if err := scanner.Scan(&r.ID, &r.ProfileID, &r.BrokerID, &r.BrokerName, &r.ResponseType, &r.EmailFrom, &r.EmailSubject, &emailBody, + &formURL, &confirmURL, &r.Confidence, &needsReviewInt, &receivedAt, &processedAt, &createdAt); err != nil { + return nil, err } - r.EmailBody = emailBody.String r.FormURL = formURL.String r.ConfirmURL = confirmURL.String r.NeedsReview = needsReviewInt == 1 - r.ReceivedAt = parseFlexibleTimeString(receivedAtStr) - r.ProcessedAt = parseFlexibleTimeString(processedAtStr) - r.CreatedAt = parseFlexibleTimeString(createdAtStr) + r.ReceivedAt = parseFlexibleTimeString(receivedAt) + r.ProcessedAt = parseFlexibleTimeString(processedAt) + r.CreatedAt = parseFlexibleTimeString(createdAt) return &r, nil } -func (s *Store) FindBrokerResponseBySubject(profileID, brokerID, subject string) (*BrokerResponse, error) { - query := `SELECT id, profile_id, broker_id, broker_name, response_type, email_from, email_subject, - form_url, confirm_url, confidence, needs_review, received_at, processed_at, created_at - FROM broker_responses WHERE profile_id = ? AND broker_id = ? AND email_subject = ? LIMIT 1` +func (s *Store) queryBrokerResponses(query string, args ...any) ([]BrokerResponse, error) { + rows, err := s.db.Query(query, args...) + if err != nil { + return nil, fmt.Errorf("failed to query broker responses: %w", err) + } + defer func() { _ = rows.Close() }() - var r BrokerResponse - var needsReviewInt int - var receivedAtStr, processedAtStr, createdAtStr sql.NullString - var formURL, confirmURL sql.NullString + var responses []BrokerResponse + for rows.Next() { + r, err := scanBrokerResponse(rows) + if err != nil { + return nil, fmt.Errorf("failed to scan broker response: %w", err) + } + responses = append(responses, *r) + } + return responses, rows.Err() +} - err := s.db.QueryRow(query, normalizeProfileID(profileID), brokerID, subject).Scan( - &r.ID, &r.ProfileID, &r.BrokerID, &r.BrokerName, &r.ResponseType, &r.EmailFrom, &r.EmailSubject, - &formURL, &confirmURL, &r.Confidence, &needsReviewInt, &receivedAtStr, &processedAtStr, &createdAtStr) - if err == sql.ErrNoRows { +func (s *Store) GetBrokerResponseByID(id int64, profileID string) (*BrokerResponse, error) { + r, err := scanBrokerResponse(s.db.QueryRow( + `SELECT `+brokerResponseCols+` FROM broker_responses WHERE id = ? AND profile_id = ?`, + id, normalizeProfileID(profileID))) + if errors.Is(err, sql.ErrNoRows) { return nil, nil } if err != nil { - return nil, fmt.Errorf("failed to find broker response: %w", err) + return nil, fmt.Errorf("failed to get broker response: %w", err) } + return r, nil +} - r.FormURL = formURL.String - r.ConfirmURL = confirmURL.String - r.NeedsReview = needsReviewInt == 1 - - return &r, nil +func (s *Store) FindBrokerResponseBySubject(profileID, brokerID, subject string) (*BrokerResponse, error) { + r, err := scanBrokerResponse(s.db.QueryRow( + `SELECT `+brokerResponseCols+` FROM broker_responses WHERE profile_id = ? AND broker_id = ? AND email_subject = ? LIMIT 1`, + normalizeProfileID(profileID), brokerID, subject)) + if errors.Is(err, sql.ErrNoRows) { + return nil, nil + } + if err != nil { + return nil, fmt.Errorf("failed to find broker response: %w", err) + } + return r, nil } // UpdateBrokerResponseClassification updates the classification fields of a @@ -844,141 +861,29 @@ func (s *Store) ClearBrokerResponses() error { // (for reclassification - a full re-scan processes the whole shared inbox // regardless of which profile is active in the caller's session) func (s *Store) GetAllBrokerResponses() ([]BrokerResponse, error) { - query := `SELECT id, profile_id, broker_id, broker_name, response_type, email_from, email_subject, email_body, - form_url, confirm_url, confidence, needs_review, received_at, processed_at, created_at - FROM broker_responses ORDER BY created_at DESC` - - rows, err := s.db.Query(query) - if err != nil { - return nil, fmt.Errorf("failed to query all broker responses: %w", err) - } - defer func() { _ = rows.Close() }() - - var responses []BrokerResponse - for rows.Next() { - var r BrokerResponse - var needsReviewInt int - var receivedAtStr, processedAtStr, createdAtStr sql.NullString - var formURL, confirmURL, emailBody sql.NullString - - err := rows.Scan(&r.ID, &r.ProfileID, &r.BrokerID, &r.BrokerName, &r.ResponseType, &r.EmailFrom, &r.EmailSubject, &emailBody, - &formURL, &confirmURL, &r.Confidence, &needsReviewInt, &receivedAtStr, &processedAtStr, &createdAtStr) - if err != nil { - return nil, fmt.Errorf("failed to scan broker response: %w", err) - } - - r.EmailBody = emailBody.String - r.FormURL = formURL.String - r.ConfirmURL = confirmURL.String - r.NeedsReview = needsReviewInt == 1 - - r.ReceivedAt = parseFlexibleTimeString(receivedAtStr) - r.ProcessedAt = parseFlexibleTimeString(processedAtStr) - r.CreatedAt = parseFlexibleTimeString(createdAtStr) - - responses = append(responses, r) - } - - return responses, rows.Err() + return s.queryBrokerResponses(`SELECT ` + brokerResponseCols + ` FROM broker_responses ORDER BY created_at DESC`) } // GetBrokerResponsesForExport returns one profile's responses, oldest first, // including email_body. Used by `eraser export`. func (s *Store) GetBrokerResponsesForExport(profileID string) ([]BrokerResponse, error) { - query := `SELECT id, profile_id, broker_id, broker_name, response_type, email_from, - email_subject, email_body, form_url, confirm_url, confidence, needs_review, - received_at, processed_at, created_at - FROM broker_responses WHERE profile_id = ? ORDER BY received_at ASC, id ASC` - - rows, err := s.db.Query(query, normalizeProfileID(profileID)) - if err != nil { - return nil, fmt.Errorf("failed to query broker responses: %w", err) - } - defer func() { _ = rows.Close() }() - - var responses []BrokerResponse - for rows.Next() { - var r BrokerResponse - var needsReviewInt int - var receivedAtStr, processedAtStr, createdAtStr sql.NullString - var formURL, confirmURL, emailBody sql.NullString - - if err := rows.Scan(&r.ID, &r.ProfileID, &r.BrokerID, &r.BrokerName, &r.ResponseType, - &r.EmailFrom, &r.EmailSubject, &emailBody, &formURL, &confirmURL, &r.Confidence, - &needsReviewInt, &receivedAtStr, &processedAtStr, &createdAtStr); err != nil { - return nil, fmt.Errorf("failed to scan broker response: %w", err) - } - r.EmailBody = emailBody.String - r.FormURL = formURL.String - r.ConfirmURL = confirmURL.String - r.NeedsReview = needsReviewInt == 1 - r.ReceivedAt = parseFlexibleTimeString(receivedAtStr) - r.ProcessedAt = parseFlexibleTimeString(processedAtStr) - r.CreatedAt = parseFlexibleTimeString(createdAtStr) - responses = append(responses, r) - } - return responses, rows.Err() + return s.queryBrokerResponses(`SELECT `+brokerResponseCols+` FROM broker_responses + WHERE profile_id = ? ORDER BY received_at ASC, id ASC`, normalizeProfileID(profileID)) } // GetBrokerResponses retrieves broker responses for one profile, with optional filtering func (s *Store) GetBrokerResponses(profileID, responseType string, needsReview bool, limit int) ([]BrokerResponse, error) { - var query string - var args []interface{} - profileID = normalizeProfileID(profileID) - - if responseType != "" && needsReview { - query = `SELECT id, profile_id, broker_id, broker_name, response_type, email_from, email_subject, - form_url, confirm_url, confidence, needs_review, received_at, processed_at, created_at - FROM broker_responses WHERE profile_id = ? AND response_type = ? AND needs_review = 1 ORDER BY created_at DESC LIMIT ?` - args = []interface{}{profileID, responseType, limit} - } else if responseType != "" { - query = `SELECT id, profile_id, broker_id, broker_name, response_type, email_from, email_subject, - form_url, confirm_url, confidence, needs_review, received_at, processed_at, created_at - FROM broker_responses WHERE profile_id = ? AND response_type = ? ORDER BY created_at DESC LIMIT ?` - args = []interface{}{profileID, responseType, limit} - } else if needsReview { - query = `SELECT id, profile_id, broker_id, broker_name, response_type, email_from, email_subject, - form_url, confirm_url, confidence, needs_review, received_at, processed_at, created_at - FROM broker_responses WHERE profile_id = ? AND needs_review = 1 ORDER BY created_at DESC LIMIT ?` - args = []interface{}{profileID, limit} - } else { - query = `SELECT id, profile_id, broker_id, broker_name, response_type, email_from, email_subject, - form_url, confirm_url, confidence, needs_review, received_at, processed_at, created_at - FROM broker_responses WHERE profile_id = ? ORDER BY created_at DESC LIMIT ?` - args = []interface{}{profileID, limit} - } - - rows, err := s.db.Query(query, args...) - if err != nil { - return nil, fmt.Errorf("failed to query broker responses: %w", err) + where := "profile_id = ?" + args := []any{normalizeProfileID(profileID)} + if responseType != "" { + where += " AND response_type = ?" + args = append(args, responseType) } - defer func() { _ = rows.Close() }() - - var responses []BrokerResponse - for rows.Next() { - var r BrokerResponse - var needsReviewInt int - var receivedAtStr, processedAtStr, createdAtStr sql.NullString - var formURL, confirmURL sql.NullString - - err := rows.Scan(&r.ID, &r.ProfileID, &r.BrokerID, &r.BrokerName, &r.ResponseType, &r.EmailFrom, &r.EmailSubject, - &formURL, &confirmURL, &r.Confidence, &needsReviewInt, &receivedAtStr, &processedAtStr, &createdAtStr) - if err != nil { - return nil, fmt.Errorf("failed to scan broker response: %w", err) - } - - r.FormURL = formURL.String - r.ConfirmURL = confirmURL.String - r.NeedsReview = needsReviewInt == 1 - - r.ReceivedAt = parseFlexibleTimeString(receivedAtStr) - r.ProcessedAt = parseFlexibleTimeString(processedAtStr) - r.CreatedAt = parseFlexibleTimeString(createdAtStr) - - responses = append(responses, r) + if needsReview { + where += " AND needs_review = 1" } - - return responses, rows.Err() + args = append(args, limit) + return s.queryBrokerResponses(`SELECT `+brokerResponseColsNoBody+` FROM broker_responses WHERE `+where+` ORDER BY created_at DESC LIMIT ?`, args...) } // GetResponseStats returns counts of response types for one profile diff --git a/internal/inbox/parser.go b/internal/inbox/parser.go index f17508c..c9b18df 100644 --- a/internal/inbox/parser.go +++ b/internal/inbox/parser.go @@ -6,7 +6,7 @@ import ( "regexp" "strings" - "github.com/PuerkitoBio/goquery" + "golang.org/x/net/html" ) // ExtractedURLs contains categorized URLs from an email @@ -159,25 +159,33 @@ func extractURLsFromText(text string) []string { // extractURLsFromHTML extracts href values. The size bound is at the MIME read // in monitor.go; html is already in memory here. -func extractURLsFromHTML(html string) []string { +func extractURLsFromHTML(doc string) []string { var urls []string - - doc, err := goquery.NewDocumentFromReader(strings.NewReader(html)) - if err != nil { - // Fallback to regex - return extractURLsFromText(html) - } - - doc.Find("a[href]").Each(func(i int, s *goquery.Selection) { - if href, exists := s.Attr("href"); exists { - urls = append(urls, href) + var text strings.Builder + z := html.NewTokenizer(strings.NewReader(doc)) + for { + switch z.Next() { + case html.ErrorToken: + // EOF, or unparseable input - either way, also check the text + // seen so far for bare URLs. + return append(urls, extractURLsFromText(text.String())...) + case html.TextToken: + text.Write(z.Text()) + text.WriteByte(' ') + case html.StartTagToken, html.SelfClosingTagToken: + name, hasAttr := z.TagName() + if string(name) != "a" { + continue + } + for hasAttr { + var key, val []byte + key, val, hasAttr = z.TagAttr() + if string(key) == "href" { + urls = append(urls, string(val)) + } + } } - }) - - // Also check for URLs in plain text within the HTML - urls = append(urls, extractURLsFromText(doc.Text())...) - - return urls + } } // cleanURL normalizes and validates a URL diff --git a/internal/inbox/parser_test.go b/internal/inbox/parser_test.go index 3a4b8d8..c866234 100644 --- a/internal/inbox/parser_test.go +++ b/internal/inbox/parser_test.go @@ -76,3 +76,18 @@ func TestIsPrivateOrLoopbackHost(t *testing.T) { }) } } + +func TestExtractURLsFromHTML(t *testing.T) { + doc := `

Confirm here + or visit https://acme.example/portal.

no href` + got := extractURLsFromHTML(doc) + want := []string{"https://acme.example/confirm?t=1&u=2", "https://acme.example/portal."} + if len(got) != len(want) { + t.Fatalf("got %q, want %q", got, want) + } + for i := range want { + if got[i] != want[i] { + t.Fatalf("got %q, want %q", got, want) + } + } +} diff --git a/internal/inbox/scan.go b/internal/inbox/scan.go new file mode 100644 index 0000000..c9de7fa --- /dev/null +++ b/internal/inbox/scan.go @@ -0,0 +1,163 @@ +package inbox + +import ( + "context" + "fmt" + "log" + + "github.com/drumandbytes/eraser/internal/history" +) + +// ScanOptions picks what ScanAndStore reads and how it treats replies it +// has seen before. +type ScanOptions struct { + Days int + // IncludeArchive also reads the archive folder, where earlier scans + // moved replies (the web UI's scans; `eraser monitor` reads INBOX only). + IncludeArchive bool + // Reclassify re-runs the classifier on replies already stored and + // updates them in place (the web UI's rescan). + Reclassify bool +} + +// ScanResult is what one ScanAndStore found. +type ScanResult struct { + Summary Summary // every broker reply found, by type + New []ClassifiedResponse // replies stored for the first time + Updated int // stored replies reclassified (Reclassify only) + Archived int +} + +// ScanAndStore fetches broker replies from a connected inbox, classifies +// them and stores new ones in the history, attributing each to the profile +// that emailed that broker (a shared inbox serves several profiles). New +// replies also advance the broker's pipeline status; replies seen before +// don't, so a re-scan can't undo progress made since. With AutoArchive, the +// INBOX replies are moved to the archive folder afterwards. +func (m *Monitor) ScanAndStore(ctx context.Context, store *history.Store, opt ScanOptions) (ScanResult, error) { + var res ScanResult + emails, err := m.FetchBrokerEmails(ctx, opt.Days) + if err != nil { + return res, fmt.Errorf("failed to fetch emails from %s: %w", m.config.Email, err) + } + // Only INBOX UIDs are archived: ArchiveEmails applies them to INBOX, + // where an archive-folder UID would name some other message. + var inboxUIDs []uint32 + for _, e := range emails { + if e.UID > 0 { + inboxUIDs = append(inboxUIDs, e.UID) + } + } + if opt.IncludeArchive && m.config.ArchiveFolder != "" { + archived, err := m.FetchBrokerEmailsFromFolder(ctx, m.config.ArchiveFolder, opt.Days) + if err != nil { + log.Printf("Warning: failed to fetch from archive folder %s: %v", m.config.ArchiveFolder, err) + } + emails = append(emails, archived...) + } + + var all []ClassifiedResponse + for i := range emails { + classified, outcome, err := RecordReply(store, &emails[i], opt.Reclassify) + if err != nil { + log.Printf("Warning: failed to store broker response for %s: %v", emails[i].BrokerID, err) + } + all = append(all, classified) + switch outcome { + case ReplyNew: + res.New = append(res.New, classified) + case ReplyUpdated: + res.Updated++ + } + } + res.Summary = SummarizeResponses(all) + + if m.config.AutoArchive && len(inboxUIDs) > 0 { + if err := m.EnsureFolderExists(m.config.ArchiveFolder); err != nil { + log.Printf("Warning: could not create archive folder %s: %v", m.config.ArchiveFolder, err) + } else if err := m.ArchiveEmails(inboxUIDs, m.config.ArchiveFolder); err != nil { + log.Printf("Warning: could not archive emails: %v", err) + } else { + res.Archived = len(inboxUIDs) + } + } + return res, nil +} + +// ReplyOutcome is what RecordReply did with a reply. +type ReplyOutcome int + +const ( + ReplySeen ReplyOutcome = iota // already stored, left as is + ReplyNew // stored for the first time + ReplyUpdated // already stored, reclassified +) + +// RecordReply classifies one broker reply and stores it (see ScanAndStore). +func RecordReply(store *history.Store, e *Email, reclassify bool) (ClassifiedResponse, ReplyOutcome, error) { + classified := ClassifyResponse(e) + body := e.Body + if body == "" { + body = e.HTMLBody + } + + profileID, err := store.ResolveProfileForBroker(e.BrokerID) + if err != nil { + profileID = history.DefaultProfileID + } + + if reclassify { + existing, err := store.FindBrokerResponseBySubject(profileID, e.BrokerID, e.Subject) + if err != nil { + return classified, ReplySeen, err + } + if existing != nil { + if err := store.UpdateBrokerResponseClassification(existing.ID, existing.ProfileID, string(classified.Type), + classified.FormURL, classified.ConfirmURL, classified.Confidence, classified.NeedsReview); err != nil { + return classified, ReplySeen, err + } + if existing.EmailBody == "" && body != "" { + if err := store.UpdateBrokerResponseBody(existing.ID, existing.ProfileID, body); err != nil { + return classified, ReplyUpdated, err + } + } + return classified, ReplyUpdated, nil + } + } + + inserted, err := store.AddBrokerResponseIfNew(&history.BrokerResponse{ + ProfileID: profileID, + BrokerID: e.BrokerID, + BrokerName: e.BrokerName, + ResponseType: string(classified.Type), + EmailFrom: e.From, + EmailSubject: e.Subject, + EmailBody: body, + FormURL: classified.FormURL, + ConfirmURL: classified.ConfirmURL, + Confidence: classified.Confidence, + NeedsReview: classified.NeedsReview, + ReceivedAt: e.ReceivedAt, + }) + if err != nil || !inserted { + return classified, ReplySeen, err + } + // No matching request row is fine: the reply still counts. + _ = store.UpdatePipelineStatus(profileID, e.BrokerID, pipelineStatusFor(classified.Type)) + return classified, ReplyNew, nil +} + +func pipelineStatusFor(t ResponseType) history.PipelineStatus { + switch t { + case ResponseSuccess: + return history.PipelineConfirmed + case ResponseFormRequired: + return history.PipelineFormRequired + case ResponseConfirmationRequired: + return history.PipelineAwaitingConfirmation + case ResponseRejected: + return history.PipelineRejected + default: + return history.PipelineAwaitingResponse + } +} diff --git a/internal/inbox/scan_test.go b/internal/inbox/scan_test.go new file mode 100644 index 0000000..bd94177 --- /dev/null +++ b/internal/inbox/scan_test.go @@ -0,0 +1,47 @@ +package inbox + +import ( + "path/filepath" + "testing" + "time" + + "github.com/drumandbytes/eraser/internal/history" +) + +func TestRecordReply(t *testing.T) { + store, err := history.NewStore(filepath.Join(t.TempDir(), "history.db")) + if err != nil { + t.Fatal(err) + } + defer func() { _ = store.Close() }() + if err := store.Add(&history.Record{ProfileID: "jane", BrokerID: "acme", BrokerName: "Acme", Email: "privacy@acme.example", Template: "gdpr", Status: history.StatusSent, SentAt: time.Now()}); err != nil { + t.Fatal(err) + } + + reply := &Email{ + From: "privacy@acme.example", BrokerID: "acme", BrokerName: "Acme", + Subject: "Re: Erasure request", + Body: "We have deleted your personal data from our systems.", + ReceivedAt: time.Date(2026, 9, 1, 10, 0, 0, 0, time.UTC), + } + + _, outcome, err := RecordReply(store, reply, false) + if err != nil || outcome != ReplyNew { + t.Fatalf("first record: outcome %v, err %v; want ReplyNew", outcome, err) + } + reqs, err := store.GetAllRequests("jane") + if err != nil || len(reqs) != 1 || reqs[0].PipelineStatus != history.PipelineConfirmed { + t.Fatalf("pipeline status not advanced by a new reply: %+v (%v)", reqs, err) + } + + if _, outcome, _ := RecordReply(store, reply, false); outcome != ReplySeen { + t.Fatalf("second scan: outcome %v, want ReplySeen", outcome) + } + if _, outcome, _ := RecordReply(store, reply, true); outcome != ReplyUpdated { + t.Fatalf("reclassify: outcome %v, want ReplyUpdated", outcome) + } + all, _ := store.GetAllBrokerResponses() + if len(all) != 1 || all[0].ProfileID != "jane" || all[0].EmailBody == "" { + t.Fatalf("want one stored reply attributed to jane with its body, got %+v", all) + } +} diff --git a/internal/web/handlers_api.go b/internal/web/handlers_api.go index 6ddb770..b3f7526 100644 --- a/internal/web/handlers_api.go +++ b/internal/web/handlers_api.go @@ -209,326 +209,71 @@ func (s *Server) handleAPIResponseReviewed(w http.ResponseWriter, r *http.Reques } func (s *Server) handleAPIInboxScan(w http.ResponseWriter, r *http.Request) { - cfg := s.getConfig() - if cfg == nil || !cfg.InboxForProfile(s.activeProfile(r)).Enabled { - _, _ = w.Write([]byte(` -
- Inbox monitoring not configured. -

Go to Settings to configure IMAP access.

-
- `)) - return - } - // Scans only the active profile's own inbox (its mail.inbox override, or - // the shared inbox: block if it doesn't have one) - `eraser monitor` - // covers every configured inbox in one run for the automated/background - // path. - inboxCfg := cfg.InboxForProfile(s.activeProfile(r)) - - monitor := inbox.NewMonitor(inboxCfg, s.brokerDB.Brokers) - - ctx, cancel := context.WithTimeout(r.Context(), 60*time.Second) - defer cancel() - - if err := monitor.Connect(ctx); err != nil { - _, _ = fmt.Fprintf(w, ` -
- Failed to connect to inbox: %s -
- `, template.HTMLEscapeString(err.Error())) - return - } - defer func() { _ = monitor.Disconnect() }() - - // Fetch emails from last 7 days - check both INBOX and archive folder - emails, err := monitor.FetchBrokerEmails(ctx, 7) - if err != nil { - _, _ = fmt.Fprintf(w, ` -
- Failed to fetch emails: %s -
- `, template.HTMLEscapeString(err.Error())) - return - } - - // Also check archive folder if configured - if inboxCfg.ArchiveFolder != "" { - archiveEmails, err := monitor.FetchBrokerEmailsFromFolder(ctx, inboxCfg.ArchiveFolder, 7) - if err != nil { - log.Printf("Warning: failed to fetch from archive folder %s: %v", inboxCfg.ArchiveFolder, err) - } else { - emails = append(emails, archiveEmails...) - } - } - - if len(emails) == 0 { - _, _ = w.Write([]byte(` -
- No new broker emails found. -

No emails from known data brokers in the last 7 days.

-
- `)) - return - } - - // Classify and store each email - var success, formRequired, confirmRequired, rejected, unknown int - var processedUIDs []uint32 // Track UIDs for archiving - for _, email := range emails { - classified := inbox.ClassifyResponse(&email) - processedUIDs = append(processedUIDs, email.UID) - - // Get body content (prefer plain text, fall back to HTML) - bodyContent := email.Body - if bodyContent == "" { - bodyContent = email.HTMLBody - } - - // A shared inbox carries replies for every profile's sent requests - // together, so attribute this reply to whichever profile actually - // emailed this broker rather than to whatever profile is "active" - // in the session running this scan. - profileID := history.DefaultProfileID - if s.historyStore != nil { - if resolved, err := s.historyStore.ResolveProfileForBroker(email.BrokerID); err == nil { - profileID = resolved - } - } - - brokerResp := &history.BrokerResponse{ - ProfileID: profileID, - BrokerID: email.BrokerID, - BrokerName: email.BrokerName, - ResponseType: string(classified.Type), - EmailFrom: email.From, - EmailSubject: email.Subject, - EmailBody: bodyContent, - FormURL: classified.FormURL, - ConfirmURL: classified.ConfirmURL, - Confidence: classified.Confidence, - NeedsReview: classified.NeedsReview, - ReceivedAt: email.ReceivedAt, - } - - if s.historyStore != nil { - if _, err := s.historyStore.AddBrokerResponseIfNew(brokerResp); err != nil { - // Don't let a DB write failure silently vanish while the - // in-memory counters below still report success - this is - // exactly what caused digisamroc/eraser#17 (Pipeline page - // empty despite "Scan Complete!" reporting matches). - log.Printf("Warning: failed to store broker response for %s: %v", brokerResp.BrokerID, err) - } - } - - // Count by type - switch classified.Type { - case inbox.ResponseSuccess: - success++ - case inbox.ResponseFormRequired: - formRequired++ - case inbox.ResponseConfirmationRequired: - confirmRequired++ - case inbox.ResponseRejected: - rejected++ - default: - unknown++ - } - } + s.scanInbox(w, r, inbox.ScanOptions{Days: 7, IncludeArchive: true}, 60*time.Second) +} - // Auto-archive processed emails to the Eraser folder - var archived int - if inboxCfg.AutoArchive && len(processedUIDs) > 0 { - if err := monitor.ArchiveEmails(processedUIDs, inboxCfg.ArchiveFolder); err != nil { - log.Printf("Warning: failed to archive emails: %v", err) - } else { - archived = len(processedUIDs) - log.Printf("Archived %d emails to %s folder", archived, inboxCfg.ArchiveFolder) +// handleAPIInboxRescan re-reads 30 days and reclassifies replies already +// stored; ?clear=true drops every stored reply first. +func (s *Server) handleAPIInboxRescan(w http.ResponseWriter, r *http.Request) { + if r.URL.Query().Get("clear") == "true" && s.historyStore != nil && s.inboxConfigured(r) { + if err := s.historyStore.ClearBrokerResponses(); err != nil { + writeScanAlert(w, "error", "Failed to clear responses:", err.Error()) + return } } - - _, _ = fmt.Fprintf(w, ` -
- Scan complete! Found %d broker emails. -
-
Success: %d
-
Form required: %d
-
Confirm required: %d
-
Rejected: %d
-
Unknown: %d
-
-

- View pending tasks | - Refresh page -

-
- `, len(emails), success, formRequired, confirmRequired, rejected, unknown) + s.scanInbox(w, r, inbox.ScanOptions{Days: 30, IncludeArchive: true, Reclassify: true}, 180*time.Second) } -// handleAPIInboxRescan rescans all emails and reclassifies them with the improved classifier -func (s *Server) handleAPIInboxRescan(w http.ResponseWriter, r *http.Request) { +func (s *Server) inboxConfigured(r *http.Request) bool { cfg := s.getConfig() - if cfg == nil || !cfg.InboxForProfile(s.activeProfile(r)).Enabled { - _, _ = w.Write([]byte(` -
- Inbox monitoring not configured. -

Go to Settings to configure IMAP access.

-
- `)) + return cfg != nil && cfg.InboxForProfile(s.activeProfile(r)).Enabled +} + +// scanInbox scans the active profile's own inbox (its mail.inbox override, +// or the shared inbox: block) - `eraser monitor` and automated cycles cover +// every configured inbox. +func (s *Server) scanInbox(w http.ResponseWriter, r *http.Request, opt inbox.ScanOptions, timeout time.Duration) { + if !s.inboxConfigured(r) { + _, _ = w.Write([]byte(`
Inbox monitoring not configured. + Go to Settings to configure IMAP access.
`)) return } - inboxCfg := cfg.InboxForProfile(s.activeProfile(r)) - - clearFirst := r.URL.Query().Get("clear") == "true" - if clearFirst && s.historyStore != nil { - if err := s.historyStore.ClearBrokerResponses(); err != nil { - _, _ = fmt.Fprintf(w, ` -
- Failed to clear responses: %s -
- `, template.HTMLEscapeString(err.Error())) - return - } + if s.historyStore == nil { + writeScanAlert(w, "error", "Database not available.", "") + return } - + inboxCfg := s.getConfig().InboxForProfile(s.activeProfile(r)) monitor := inbox.NewMonitor(inboxCfg, s.brokerDB.Brokers) - ctx, cancel := context.WithTimeout(r.Context(), 180*time.Second) + ctx, cancel := context.WithTimeout(r.Context(), timeout) defer cancel() - if err := monitor.Connect(ctx); err != nil { - _, _ = fmt.Fprintf(w, ` -
- Failed to connect to inbox: %s -
- `, template.HTMLEscapeString(err.Error())) + writeScanAlert(w, "error", "Failed to connect to inbox:", err.Error()) return } defer func() { _ = monitor.Disconnect() }() - // Fetch emails from last 30 days for full rescan - check both INBOX and archive folder - emails, err := monitor.FetchBrokerEmails(ctx, 30) + res, err := monitor.ScanAndStore(ctx, s.historyStore, opt) if err != nil { - _, _ = fmt.Fprintf(w, ` -
- Failed to fetch emails: %s -
- `, template.HTMLEscapeString(err.Error())) + writeScanAlert(w, "error", "Failed to fetch emails:", err.Error()) return } - - // Also check archive folder if configured - if inboxCfg.ArchiveFolder != "" { - archiveEmails, err := monitor.FetchBrokerEmailsFromFolder(ctx, inboxCfg.ArchiveFolder, 30) - if err != nil { - log.Printf("Warning: failed to fetch from archive folder %s: %v", inboxCfg.ArchiveFolder, err) - } else { - emails = append(emails, archiveEmails...) - } - } - - if len(emails) == 0 { - _, _ = w.Write([]byte(` -
- No broker emails found. -

No emails from known data brokers in the last 30 days.

-
- `)) + if res.Summary.Total == 0 { + writeScanAlert(w, "info", "No broker emails found.", fmt.Sprintf("No emails from known data brokers in the last %d days.", opt.Days)) return } - // Classify and store/update each email - var success, formRequired, confirmRequired, rejected, pending, unknown int - var updated, inserted int - for _, email := range emails { - classified := inbox.ClassifyResponse(&email) - - // Get body content (prefer plain text, fall back to HTML) - bodyContent := email.Body - if bodyContent == "" { - bodyContent = email.HTMLBody - } - - // A shared inbox carries replies for every profile's sent requests - // together, so attribute this reply to whichever profile actually - // emailed this broker rather than to whatever profile is "active" - // in the session running this rescan. - profileID := history.DefaultProfileID - if s.historyStore != nil { - if resolved, err := s.historyStore.ResolveProfileForBroker(email.BrokerID); err == nil { - profileID = resolved - } - } - - if s.historyStore != nil { - existing, _ := s.historyStore.FindBrokerResponseBySubject(profileID, email.BrokerID, email.Subject) - if existing != nil { - err := s.historyStore.UpdateBrokerResponseClassification( - existing.ID, - existing.ProfileID, - string(classified.Type), - classified.FormURL, - classified.ConfirmURL, - classified.Confidence, - classified.NeedsReview, - ) - if err == nil { - updated++ - } else { - log.Printf("Warning: failed to update broker response classification for %s: %v", email.BrokerID, err) - } - // Also update the body if it was empty - if existing.EmailBody == "" && bodyContent != "" { - if err := s.historyStore.UpdateBrokerResponseBody(existing.ID, existing.ProfileID, bodyContent); err != nil { - log.Printf("Warning: failed to update broker response body for %s: %v", email.BrokerID, err) - } - } - } else { - // Insert new response - brokerResp := &history.BrokerResponse{ - ProfileID: profileID, - BrokerID: email.BrokerID, - BrokerName: email.BrokerName, - ResponseType: string(classified.Type), - EmailFrom: email.From, - EmailSubject: email.Subject, - EmailBody: bodyContent, - FormURL: classified.FormURL, - ConfirmURL: classified.ConfirmURL, - Confidence: classified.Confidence, - NeedsReview: classified.NeedsReview, - ReceivedAt: email.ReceivedAt, - } - if err := s.historyStore.AddBrokerResponse(brokerResp); err == nil { - inserted++ - } else { - log.Printf("Warning: failed to store broker response for %s: %v", email.BrokerID, err) - } - } - } - - // Count by type - switch classified.Type { - case inbox.ResponseSuccess: - success++ - case inbox.ResponseFormRequired: - formRequired++ - case inbox.ResponseConfirmationRequired: - confirmRequired++ - case inbox.ResponseRejected: - rejected++ - case inbox.ResponsePending: - pending++ - default: - unknown++ - } + sum := res.Summary + updated := "" + if opt.Reclassify { + updated = fmt.Sprintf(`
Updated: %d
`, res.Updated) } - _, _ = fmt.Fprintf(w, ` -
- Rescan complete! Processed %d broker emails. +
+ Scan complete! Found %d broker emails.
-
Updated: %d
New: %d
+ %s
Success: %d
@@ -543,36 +288,29 @@ func (s *Server) handleAPIInboxRescan(w http.ResponseWriter, r *http.Request) { Refresh page

- `, len(emails), updated, inserted, success, formRequired, confirmRequired, pending, rejected, unknown) + `, sum.Total, len(res.New), updated, sum.Success, sum.FormRequired, sum.ConfirmRequired, sum.Pending, sum.Rejected, sum.Unknown) +} + +// writeScanAlert writes an .alert- banner (layout.html); detail is escaped. +func writeScanAlert(w http.ResponseWriter, kind, title, detail string) { + _, _ = fmt.Fprintf(w, `
%s %s
`, kind, title, template.HTMLEscapeString(detail)) } // handleAPIReclassify reclassifies all existing database records using subject-only patterns func (s *Server) handleAPIReclassify(w http.ResponseWriter, r *http.Request) { if s.historyStore == nil { - _, _ = w.Write([]byte(` -
- Database not available. -
- `)) + writeScanAlert(w, "error", "Database not available.", "") return } responses, err := s.historyStore.GetAllBrokerResponses() if err != nil { - _, _ = fmt.Fprintf(w, ` -
- Failed to get responses: %s -
- `, template.HTMLEscapeString(err.Error())) + writeScanAlert(w, "error", "Failed to get responses:", err.Error()) return } if len(responses) == 0 { - _, _ = w.Write([]byte(` -
- No responses to reclassify. -
- `)) + writeScanAlert(w, "info", "No responses to reclassify.", "") return } @@ -729,7 +467,7 @@ func (s *Server) handleAPIReclassify(w http.ResponseWriter, r *http.Request) { } _, _ = fmt.Fprintf(w, ` -
+
Reclassification complete! Processed %d records.
Updated: %d
diff --git a/internal/web/handlers_jobs.go b/internal/web/handlers_jobs.go index c82d2cf..4eba688 100644 --- a/internal/web/handlers_jobs.go +++ b/internal/web/handlers_jobs.go @@ -10,127 +10,17 @@ import ( "strings" "time" - "github.com/drumandbytes/eraser/internal/broker" "github.com/drumandbytes/eraser/internal/config" "github.com/drumandbytes/eraser/internal/email" "github.com/drumandbytes/eraser/internal/history" "github.com/go-chi/chi/v5" ) -// checkPendingJob resumes each profile's incomplete job; state is per -// profile since two profiles can send concurrently. -func (s *Server) checkPendingJob() { - cfg := s.getConfig() - if cfg == nil { - return - } - for _, p := range cfg.GetProfiles() { - s.checkPendingJobForProfile(p.ID) - } -} - -func (s *Server) checkPendingJobForProfile(profileID string) { - state, err := s.jobPersistence.Load(profileID) - if err != nil { - log.Printf("Warning: failed to load pending job for profile %s: %v", profileID, err) - return - } - - if state == nil || len(state.RemainingBrokers) == 0 { - return // No pending job - } - - fmt.Printf("\nFound incomplete send job for profile %s: %d of %d brokers remaining\n", profileID, len(state.RemainingBrokers), state.Total) - fmt.Printf("Already sent: %d, failed: %d\n", state.Sent, state.Failed) - - // Auto-resume the job - go s.resumePendingJob(state) -} - -// resumePendingJob resumes processing of an incomplete job -func (s *Server) resumePendingJob(state *PersistentJobState) { - // Wait a moment for the server to fully start - time.Sleep(2 * time.Second) - - // jobs from before multi-profile have no ProfileID; resolve it up front so - // every Clear below targets the file Load read - profileID := state.ProfileID - if profileID == "" { - profileID = config.DefaultProfileID - } - - cfg := s.getConfig() - if cfg == nil { - log.Printf("Cannot resume job: email not configured") - _ = s.jobPersistence.Clear(profileID) - return - } - - activeProfile, err := cfg.GetProfile(profileID) - if err != nil { - if profiles := cfg.GetProfiles(); len(profiles) > 0 { - activeProfile = profiles[0] - } - } - - emailCfg := cfg.EmailForProfile(activeProfile) - if emailCfg.Provider == "" { - log.Printf("Cannot resume job: email not configured") - _ = s.jobPersistence.Clear(profileID) - return - } - - sender, err := email.NewSender(emailCfg) - if err != nil { - log.Printf("Cannot resume job: failed to create email sender: %v", err) - _ = s.jobPersistence.Clear(profileID) - return - } - - brokerMap := make(map[string]broker.Broker) - for _, b := range s.brokerDB.Brokers { - brokerMap[b.ID] = b - } - - // Anything sent since the job paused (an automatic run, the CLI, a - // "Send all" click) is no longer due; resending it would double-email. - var statuses map[string]history.BrokerStatus - if s.historyStore != nil { - statuses, _ = s.historyStore.GetAllBrokerStatuses(profileID) - } - var toSend []BrokerWithStatus - for _, id := range state.RemainingBrokers { - b, ok := brokerMap[id] - if !ok { - continue - } - if st, sent := statuses[id]; sent && st.Status == history.StatusSent && time.Since(st.LastSent) < history.ResendCooldown { - continue - } - toSend = append(toSend, BrokerWithStatus{Broker: b, Status: "never"}) - } - - if len(toSend) == 0 { - log.Printf("No valid brokers remaining in pending job") - _ = s.jobPersistence.Clear(profileID) - return - } - - // Create a new job to continue processing, preserving the profile the - // original job was scoped to. - job := s.jobManager.Create(state.Total, profileID) - job.Update(state.Sent, state.Failed, "", "") - - fmt.Printf("Resuming send job: %d brokers remaining...\n", len(toSend)) - - s.processSendJob(job, toSend, sender) -} - func (s *Server) handleAPISendOne(w http.ResponseWriter, r *http.Request) { // Rate limiting - prevent abuse of email sending if !s.rateLimiter.Allow("send") { w.WriteHeader(http.StatusTooManyRequests) - _, _ = w.Write([]byte(`Rate limit exceeded. Please wait a moment before sending more emails.`)) + _, _ = w.Write([]byte(`Rate limit exceeded. Please wait a moment before sending more emails.`)) return } @@ -139,14 +29,14 @@ func (s *Server) handleAPISendOne(w http.ResponseWriter, r *http.Request) { br := s.brokerDB.FindByID(brokerID) if br == nil { w.WriteHeader(http.StatusNotFound) - _, _ = w.Write([]byte(`Broker not found`)) + _, _ = w.Write([]byte(`Broker not found`)) return } cfg := s.getConfig() if cfg == nil { w.WriteHeader(http.StatusBadRequest) - _, _ = w.Write([]byte(`Email not configured. Configure now`)) + _, _ = w.Write([]byte(`Email not configured. Configure now`)) return } @@ -154,81 +44,50 @@ func (s *Server) handleAPISendOne(w http.ResponseWriter, r *http.Request) { emailCfg := cfg.EmailForProfile(activeProfile) if emailCfg.Provider == "" { w.WriteHeader(http.StatusBadRequest) - _, _ = w.Write([]byte(`Email not configured. Configure now`)) + _, _ = w.Write([]byte(`Email not configured. Configure now`)) return } if cfg.Options.DryRun { w.WriteHeader(http.StatusBadRequest) - _, _ = w.Write([]byte(`Web sending is disabled while options.dry_run is true.`)) + _, _ = w.Write([]byte(`Web sending is disabled while options.dry_run is true.`)) return } if br.Email == "" { w.WriteHeader(http.StatusBadRequest) - _, _ = w.Write([]byte(`No email on file - needs manual follow-up (check for an opt-out form/portal)`)) + _, _ = w.Write([]byte(`No email on file - needs manual follow-up (check for an opt-out form/portal)`)) return } sender, err := email.NewSender(emailCfg) if err != nil { - _, _ = fmt.Fprintf(w, `Error: %s`, template.HTMLEscapeString(err.Error())) + _, _ = fmt.Fprintf(w, `Error: %s`, template.HTMLEscapeString(err.Error())) return } - // configured template (gdpr/ccpa/generic); config.Load guarantees one - tmplName := cfg.Options.Template - rendered, err := s.tmplEngine.Render(tmplName, activeProfile.Profile, *br) - if err != nil { - _, _ = fmt.Fprintf(w, `Template error: %s`, template.HTMLEscapeString(err.Error())) - return - } - - msg := email.Message{ - To: br.Email, - From: emailCfg.From, - Subject: rendered.Subject, - Body: rendered.Body, - } - ctx, cancel := context.WithTimeout(r.Context(), 30*time.Second) defer cancel() - - result := sender.Send(ctx, msg) - - // Record in history - record := &history.Record{ - ProfileID: activeProfile.ID, - BrokerID: br.ID, - BrokerName: br.Name, - Email: br.Email, - Template: tmplName, - SentAt: time.Now(), + record, err := email.SendRemoval(ctx, sender, s.tmplEngine, cfg.Options.Template, activeProfile, emailCfg.From, *br) + if err != nil { + _, _ = fmt.Fprintf(w, `Template error: %s`, template.HTMLEscapeString(err.Error())) + return } + s.recordSend(record) - if result.Success { - record.Status = history.StatusSent - record.MessageID = result.MessageID + if record.Status == history.StatusSent { + _, _ = w.Write([]byte(`Sent`)) } else { - record.Status = history.StatusFailed - if result.Error != nil { - record.Error = result.Error.Error() - } + _, _ = fmt.Fprintf(w, `Failed`, template.HTMLEscapeString(record.Error)) } +} - if s.historyStore != nil { - if err := s.historyStore.Add(record); err != nil { - log.Printf("Warning: failed to record send to %s in history: %v", record.BrokerID, err) - } +// recordSend stores a send attempt, logging (not failing on) a write error. +func (s *Server) recordSend(record *history.Record) { + if s.historyStore == nil { + return } - - if result.Success { - _, _ = w.Write([]byte(`Sent`)) - } else { - errMsg := "Unknown error" - if result.Error != nil { - errMsg = result.Error.Error() - } - _, _ = fmt.Fprintf(w, `Failed`, template.HTMLEscapeString(errMsg)) + if err := s.historyStore.Add(record); err != nil { + log.Printf("Warning: failed to record send to %s in history: %v", record.BrokerID, err) } } @@ -353,29 +212,6 @@ func (s *Server) handleAPISendAll(w http.ResponseWriter, r *http.Request) { return } - brokerIDs := make([]string, len(toSend)) - for i, b := range toSend { - brokerIDs[i] = b.ID - } - - jobState := &PersistentJobState{ - ID: job.ID, - ProfileID: activeProfile.ID, - Status: job.GetStatus(), - Sent: 0, - Failed: 0, - Total: len(toSend), - StartedAt: job.StartedAt, - RemainingBrokers: brokerIDs, - Search: search, - Category: category, - Region: region, - StatusFilter: status, - } - if err := s.jobPersistence.Save(jobState); err != nil { - log.Printf("Warning: failed to save job state: %v", err) - } - go s.processSendJob(job, toSend, sender) _ = json.NewEncoder(w).Encode(map[string]interface{}{ @@ -439,12 +275,6 @@ func (s *Server) processSendJob(job *Job, toSend []BrokerWithStatus, sender *ema } } - // Track remaining brokers for persistence - remaining := make([]string, len(toSend)) - for i, b := range toSend { - remaining[i] = b.ID - } - for i, b := range toSend { if job.IsCancelled() { break @@ -456,90 +286,38 @@ func (s *Server) processSendJob(job *Job, toSend []BrokerWithStatus, sender *ema if s.inAppScheduling() || s.osInstalled() { next = "Automation will send the rest once the limit frees up." } - job.Pause(sent, fmt.Sprintf("Daily limit of %d emails reached. %d brokers remaining. %s", dailyLimit, len(remaining), next)) - s.saveJobProgress(job, sent, failed, remaining) - log.Printf("Job paused: daily limit of %d reached (%d already sent today, %d this run), %d remaining", dailyLimit, alreadySentToday, sent, len(remaining)) + left := len(toSend) - i + job.Pause(sent, fmt.Sprintf("Daily limit of %d emails reached. %d brokers remaining. %s", dailyLimit, left, next)) + log.Printf("Job paused: daily limit of %d reached (%d already sent today, %d this run), %d remaining", dailyLimit, alreadySentToday, sent, left) return } job.Update(sent, failed, b.Name, b.ID) - // Generate email using the user's configured template (see the - // same fix/comment in handleAPISendOne above) - rendered, err := s.tmplEngine.Render(cfg.Options.Template, activeProfile.Profile, b.Broker) + ctx, cancel := context.WithTimeout(job.Context(), 30*time.Second) + record, err := email.SendRemoval(ctx, sender, s.tmplEngine, cfg.Options.Template, activeProfile, cfg.EmailForProfile(activeProfile).From, b.Broker) + cancel() if err != nil { failed++ job.Update(sent, failed, b.Name, b.ID) - remaining = remaining[1:] - s.saveJobProgress(job, sent, failed, remaining) continue } + s.recordSend(record) - msg := email.Message{ - To: b.Email, - From: cfg.EmailForProfile(activeProfile).From, - Subject: rendered.Subject, - Body: rendered.Body, - } - - // Use job's context with timeout for cancellation support - ctx, cancel := context.WithTimeout(job.Context(), 30*time.Second) - result := sender.Send(ctx, msg) - cancel() - - // 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 + if record.Status == history.StatusSent { sent++ - job.ResetAuthFailures() // Reset on success + job.ResetAuthFailures() } else { - record.Status = history.StatusFailed - errMsg := "" - if result.Error != nil { - errMsg = result.Error.Error() - record.Error = errMsg - } failed++ - - // Check for auth failures and stop if too many consecutive - if strings.Contains(strings.ToLower(errMsg), "auth") { - if job.RecordAuthFailure() { - if s.historyStore != nil { - if err := s.historyStore.Add(record); err != nil { - log.Printf("Warning: failed to record send to %s in history: %v", record.BrokerID, err) - } - } - remaining = remaining[1:] - s.saveJobProgress(job, sent, failed, remaining) - job.StopWithError("auth", "Stopped due to repeated authentication failures. Your email provider may have rate-limited or blocked your account. Please check your email settings and try again later.") - log.Printf("Job stopped: repeated auth failures after %d sent, %d failed", sent, failed) - return - } - } - } - - if s.historyStore != nil { - if err := s.historyStore.Add(record); err != nil { - log.Printf("Warning: failed to record send to %s in history: %v", record.BrokerID, err) + if strings.Contains(strings.ToLower(record.Error), "auth") && job.RecordAuthFailure() { + job.StopWithError("auth", "Stopped due to repeated authentication failures. Your email provider may have rate-limited or blocked your account. Please check your email settings and try again later.") + log.Printf("Job stopped: repeated auth failures after %d sent, %d failed", sent, failed) + return } } job.Update(sent, failed, b.Name, b.ID) - // Remove processed broker from remaining and save state - remaining = remaining[1:] - s.saveJobProgress(job, sent, failed, remaining) - // Rate limit delay (skip on last item) if i < len(toSend)-1 && !job.IsCancelled() { time.Sleep(time.Duration(rateLimitMs) * time.Millisecond) @@ -547,26 +325,6 @@ func (s *Server) processSendJob(job *Job, toSend []BrokerWithStatus, sender *ema } job.Complete() - if err := s.jobPersistence.Clear(job.ProfileID); err != nil { - log.Printf("Warning: failed to clear job state: %v", err) - } -} - -// saveJobProgress saves the current job progress to disk -func (s *Server) saveJobProgress(job *Job, sent, failed int, remaining []string) { - state := &PersistentJobState{ - ID: job.ID, - ProfileID: job.ProfileID, - Status: job.GetStatus(), - Sent: sent, - Failed: failed, - Total: job.Total, - StartedAt: job.StartedAt, - RemainingBrokers: remaining, - } - if err := s.jobPersistence.Save(state); err != nil { - log.Printf("Warning: failed to save job progress: %v", err) - } } // handleAPIJobActive returns the currently running job (if any) diff --git a/internal/web/handlers_pages.go b/internal/web/handlers_pages.go index de73e1a..63d7934 100644 --- a/internal/web/handlers_pages.go +++ b/internal/web/handlers_pages.go @@ -358,13 +358,13 @@ func (s *Server) handleTaskComplete(w http.ResponseWriter, r *http.Request) { if s.historyStore == nil { w.WriteHeader(http.StatusInternalServerError) - _, _ = w.Write([]byte(`Database not available`)) + _, _ = w.Write([]byte(`Database not available`)) return } if err := s.historyStore.CompletePendingTask(taskID, s.activeProfile(r).ID, status); err != nil { w.WriteHeader(http.StatusInternalServerError) - _, _ = fmt.Fprintf(w, `Error: %s`, template.HTMLEscapeString(err.Error())) + _, _ = fmt.Fprintf(w, `Error: %s`, template.HTMLEscapeString(err.Error())) return } @@ -379,13 +379,13 @@ func (s *Server) handleTaskSkip(w http.ResponseWriter, r *http.Request) { if s.historyStore == nil { w.WriteHeader(http.StatusInternalServerError) - _, _ = w.Write([]byte(`Database not available`)) + _, _ = w.Write([]byte(`Database not available`)) return } if err := s.historyStore.CompletePendingTask(taskID, s.activeProfile(r).ID, "skipped"); err != nil { w.WriteHeader(http.StatusInternalServerError) - _, _ = fmt.Fprintf(w, `Error: %s`, template.HTMLEscapeString(err.Error())) + _, _ = fmt.Fprintf(w, `Error: %s`, template.HTMLEscapeString(err.Error())) return } diff --git a/internal/web/handlers_setup.go b/internal/web/handlers_setup.go index 075bae5..0c2b57d 100644 --- a/internal/web/handlers_setup.go +++ b/internal/web/handlers_setup.go @@ -181,14 +181,14 @@ func (s *Server) handleSetupTestSend(w http.ResponseWriter, r *http.Request) { if session == nil || session.Email.Provider == "" { w.WriteHeader(http.StatusBadRequest) - _, _ = w.Write([]byte(`
Email not configured. Please go back to the email step.
`)) + _, _ = w.Write([]byte(`
Email not configured. Please go back to the email step.
`)) return } sender, err := email.NewSender(session.Email) if err != nil { _, _ = fmt.Fprintf(w, ` -
+
Configuration error: %s

Please check your email settings and try again.

@@ -220,12 +220,12 @@ Eraser`, session.Profile.FirstName), errMsg = result.Error.Error() } _, _ = fmt.Fprintf(w, ` -
+
Test failed: %s

Please check your email configuration and try again.

@@ -234,12 +234,12 @@ Eraser`, session.Profile.FirstName), } _, _ = w.Write([]byte(` -
+
Success! Test email sent to your address.

Check your inbox (and spam folder) for the test message.

diff --git a/internal/web/job.go b/internal/web/job.go index f8eabf4..dfdae39 100644 --- a/internal/web/job.go +++ b/internal/web/job.go @@ -2,14 +2,10 @@ package web import ( "context" + "crypto/rand" "encoding/json" - "os" - "path/filepath" "sync" "time" - - "github.com/drumandbytes/eraser/internal/config" - "github.com/google/uuid" ) // JobStatus represents the status of a background job @@ -227,7 +223,7 @@ func (jm *JobManager) createLocked(total int, profileID string) *Job { ctx, cancel := context.WithCancel(context.Background()) job := &Job{ - ID: uuid.New().String(), + ID: rand.Text(), ProfileID: profileID, Status: JobStatusRunning, Progress: 0, @@ -306,81 +302,3 @@ func (jm *JobManager) Cleanup(maxAge time.Duration) { } } } - -// PersistentJobState represents a job that can be saved/loaded from disk -type PersistentJobState struct { - ID string `json:"id"` - ProfileID string `json:"profile_id"` - Status JobStatus `json:"status"` - Sent int `json:"sent"` - Failed int `json:"failed"` - Total int `json:"total"` - StartedAt time.Time `json:"started_at"` - RemainingBrokers []string `json:"remaining_brokers"` // Broker IDs still to process - Search string `json:"search"` // Original filter params - Category string `json:"category"` - Region string `json:"region"` - StatusFilter string `json:"status_filter"` -} - -// JobPersistence handles saving/loading job state -type JobPersistence struct { - dataDir string -} - -// NewJobPersistence creates a new job persistence handler -func NewJobPersistence(dataDir string) *JobPersistence { - return &JobPersistence{dataDir: dataDir} -} - -// filePath: one file per profile (they can send concurrently). The default -// profile keeps the old bare filename so an in-flight job survives upgrade. -func (jp *JobPersistence) filePath(profileID string) string { - if profileID == "" || profileID == config.DefaultProfileID { - return filepath.Join(jp.dataDir, "pending_job.json") - } - return filepath.Join(jp.dataDir, "pending_job-"+config.SlugifyID(profileID)+".json") -} - -// Save saves the job state to disk, keyed by state.ProfileID. -func (jp *JobPersistence) Save(state *PersistentJobState) error { - if err := os.MkdirAll(jp.dataDir, 0700); err != nil { - return err - } - - data, err := json.MarshalIndent(state, "", " ") - if err != nil { - return err - } - - return os.WriteFile(jp.filePath(state.ProfileID), data, 0600) -} - -// Load loads a pending job state for profileID from disk, returns nil if -// none exists. -func (jp *JobPersistence) Load(profileID string) (*PersistentJobState, error) { - data, err := os.ReadFile(jp.filePath(profileID)) - if os.IsNotExist(err) { - return nil, nil - } - if err != nil { - return nil, err - } - - var state PersistentJobState - if err := json.Unmarshal(data, &state); err != nil { - return nil, err - } - - return &state, nil -} - -// Clear removes the saved job state -// Clear removes profileID's saved job state. -func (jp *JobPersistence) Clear(profileID string) error { - err := os.Remove(jp.filePath(profileID)) - if os.IsNotExist(err) { - return nil - } - return err -} diff --git a/internal/web/job_test.go b/internal/web/job_test.go index 45406b8..17ea7cc 100644 --- a/internal/web/job_test.go +++ b/internal/web/job_test.go @@ -2,8 +2,6 @@ package web import ( "encoding/json" - "os" - "path/filepath" "sync" "testing" ) @@ -212,57 +210,3 @@ func TestCreateIfNoActiveIsRaceSafe(t *testing.T) { t.Fatalf("expected exactly 1 of %d concurrent CreateIfNoActive calls to win, got %d", attempts, wins) } } - -// TestJobPersistencePerProfileFiles is the collision JobPersistence exists -// to avoid: two profiles saving concurrently used to share one -// pending_job.json, so the second Save silently overwrote the first and a -// restart could only ever resume (or forget) whichever wrote last. -func TestJobPersistencePerProfileFiles(t *testing.T) { - jp := NewJobPersistence(t.TempDir()) - - stateA := &PersistentJobState{ID: "job-a", ProfileID: "profile-a", Total: 3, RemainingBrokers: []string{"x"}} - stateB := &PersistentJobState{ID: "job-b", ProfileID: "profile-b", Total: 5, RemainingBrokers: []string{"y", "z"}} - - if err := jp.Save(stateA); err != nil { - t.Fatalf("Save(profile-a): %v", err) - } - if err := jp.Save(stateB); err != nil { - t.Fatalf("Save(profile-b): %v", err) - } - - gotA, err := jp.Load("profile-a") - if err != nil || gotA == nil || gotA.ID != "job-a" { - t.Fatalf("Load(profile-a) = %+v, %v, want job-a", gotA, err) - } - gotB, err := jp.Load("profile-b") - if err != nil || gotB == nil || gotB.ID != "job-b" { - t.Fatalf("Load(profile-b) = %+v, %v, want job-b", gotB, err) - } - - // Clearing one profile's state doesn't touch the other's. - if err := jp.Clear("profile-a"); err != nil { - t.Fatalf("Clear(profile-a): %v", err) - } - if got, err := jp.Load("profile-a"); err != nil || got != nil { - t.Fatalf("Load(profile-a) after Clear = %+v, %v, want nil", got, err) - } - if got, err := jp.Load("profile-b"); err != nil || got == nil { - t.Fatalf("Load(profile-b) after clearing profile-a = %+v, %v, want job-b still present", got, err) - } -} - -// The default profile must keep the legacy pending_job.json filename across upgrade. -func TestJobPersistenceDefaultProfileUsesLegacyFilename(t *testing.T) { - dir := t.TempDir() - jp := NewJobPersistence(dir) - - if err := jp.Save(&PersistentJobState{ID: "job-legacy", ProfileID: "default", Total: 1}); err != nil { - t.Fatalf("Save: %v", err) - } - if _, err := os.Stat(filepath.Join(dir, "pending_job.json")); err != nil { - t.Fatalf("expected legacy pending_job.json to exist: %v", err) - } - if _, err := os.Stat(filepath.Join(dir, "pending_job-default.json")); !os.IsNotExist(err) { - t.Fatalf("expected no pending_job-default.json, got err=%v", err) - } -} diff --git a/internal/web/server.go b/internal/web/server.go index a80ddb4..7d71d98 100644 --- a/internal/web/server.go +++ b/internal/web/server.go @@ -40,6 +40,8 @@ const ( defaultSessionTTL = 30 * time.Minute ) +// RateLimiter caps requests per key in a sliding window. Keys are a fixed +// handful of endpoint names, so the map never needs pruning. type RateLimiter struct { mu sync.Mutex requests map[string][]time.Time @@ -53,7 +55,6 @@ func NewRateLimiter(limit int, window time.Duration) *RateLimiter { limit: limit, window: window, } - go rl.cleanupLoop() return rl } @@ -83,44 +84,24 @@ func (rl *RateLimiter) Allow(key string) bool { return true } -func (rl *RateLimiter) cleanupLoop() { - ticker := time.NewTicker(time.Minute) - defer ticker.Stop() - - for range ticker.C { - rl.mu.Lock() - windowStart := time.Now().Add(-rl.window) - for key, times := range rl.requests { - recent := rl.filterRecent(times, windowStart) - if len(recent) == 0 { - delete(rl.requests, key) - } else { - rl.requests[key] = recent - } - } - rl.mu.Unlock() - } -} - // Version is the build version shown in the web UI footer. main sets it from // its own -ldflags-injected version at startup; it stays "dev" otherwise. var Version = "dev" type Server struct { - config atomic.Pointer[config.Config] - configPath string - brokerDB *broker.BrokerDatabase - historyStore *history.Store - tmplEngine *emaTemplate.Engine - templates map[string]*template.Template - httpServer *http.Server - port int - csrfKey []byte - sessions *SessionStore - rateLimiter *RateLimiter - jobManager *JobManager - jobPersistence *JobPersistence - dataDir string // config directory: job state, schedule lock and state + config atomic.Pointer[config.Config] + configPath string + brokerDB *broker.BrokerDatabase + historyStore *history.Store + tmplEngine *emaTemplate.Engine + templates map[string]*template.Template + httpServer *http.Server + port int + csrfKey []byte + sessions *SessionStore + rateLimiter *RateLimiter + jobManager *JobManager + dataDir string // config directory: schedule lock and state // In-app scheduler (scheduler.go). The OS hooks are fields so tests // don't touch the real launchd/systemd setup. @@ -147,23 +128,31 @@ func NewServer(port int, cfg *config.Config, configPath string, brokerDB *broker } s := &Server{ - configPath: configPath, - brokerDB: brokerDB, - historyStore: historyStore, - tmplEngine: tmplEngine, - port: port, - csrfKey: csrfKey, - sessions: NewSessionStore(defaultSessionTTL), - rateLimiter: NewRateLimiter(defaultRateLimit, defaultRateWindow), - jobManager: NewJobManager(), - jobPersistence: NewJobPersistence(dataDir), - dataDir: dataDir, - osInstalled: schedule.Installed, - installOS: schedule.Install, - removeOS: schedule.Remove, + configPath: configPath, + brokerDB: brokerDB, + historyStore: historyStore, + tmplEngine: tmplEngine, + port: port, + csrfKey: csrfKey, + sessions: NewSessionStore(defaultSessionTTL), + rateLimiter: NewRateLimiter(defaultRateLimit, defaultRateWindow), + jobManager: NewJobManager(), + dataDir: dataDir, + osInstalled: schedule.Installed, + installOS: schedule.Install, + removeOS: schedule.Remove, } s.config.Store(cfg) + // Paused send jobs used to be saved here for resume-on-restart; that's + // gone, so drop any leftovers. + if configPath != "" { + leftovers, _ := filepath.Glob(filepath.Join(dataDir, "pending_job*.json")) + for _, f := range leftovers { + _ = os.Remove(f) + } + } + tmpl, err := s.parseTemplates() if err != nil { return nil, fmt.Errorf("failed to parse templates: %w", err) @@ -329,9 +318,6 @@ func (s *Server) Start() error { IdleTimeout: 60 * time.Second, } - // Check for pending job and offer to resume - s.checkPendingJob() - ctx, cancel := context.WithCancel(context.Background()) s.stopScheduler = cancel go s.runScheduler(ctx) diff --git a/internal/web/templates/brokers.html b/internal/web/templates/brokers.html index c76fbb3..6377934 100644 --- a/internal/web/templates/brokers.html +++ b/internal/web/templates/brokers.html @@ -65,7 +65,7 @@

Data Brokers

Sending {{.Filtered}} emails will take multiple days. To avoid email provider rate limits, Eraser sends up to {{.DailyLimit}} emails per day. - Remaining emails will automatically continue when you restart the app the next day. + The rest go out with the next automatic run (Settings → Automation), or click Send all again the next day.

@@ -415,8 +415,8 @@

Email Authenticati

Daily Limit Reached

-

Sent ${job.sent} emails today. Remaining emails will continue tomorrow.

-

${job.total - job.sent - job.failed} brokers will be sent automatically when you restart the app tomorrow.

+

Sent ${job.sent} emails today.

+

${job.total - job.sent - job.failed} brokers left. They go out with the next automatic run (Settings → Automation), or click Send all again tomorrow.