diff --git a/internal/app/mailboximport/runner.go b/internal/app/mailboximport/runner.go index 52e9076e0..a7ccb7129 100644 --- a/internal/app/mailboximport/runner.go +++ b/internal/app/mailboximport/runner.go @@ -8,6 +8,7 @@ import ( "html" "strings" "sync" + "sync/atomic" "time" "github.com/google/uuid" @@ -839,30 +840,44 @@ func (s *Service) classifyHosts(ctx context.Context) { } // publish tells the workspace an import moved, at most once a second per import unless final; -// an update inside the second is sent at its end. +// an update inside the second is sent at its end with the latest status, and a final one cancels it. func (s *Service) publish(ctx context.Context, orgID, importID uuid.UUID, status string, force bool) { if s.publisher == nil { return } now := time.Now() - if !force { - if last, ok := s.progress.Load(importID); ok { - if wait := time.Second - now.Sub(last.(time.Time)); wait > 0 { - // Deferred, not dropped: the last update of a burst is the one a watcher needs. - if _, pending := s.trailing.LoadOrStore(importID, true); !pending { - time.AfterFunc(wait, func() { - s.trailing.Delete(importID) - s.publish(context.WithoutCancel(ctx), orgID, importID, status, true) - }) + if force { + if v, ok := s.trailing.LoadAndDelete(importID); ok { + v.(*deferredPublish).timer.Stop() + } + } else if last, ok := s.progress.Load(importID); ok { + if wait := time.Second - now.Sub(last.(time.Time)); wait > 0 { + d := &deferredPublish{} + d.status.Store(status) + // The timer exists before d is shared, so a final publish can always stop it; an + // unregistered d fires into a CompareAndDelete that fails and does nothing. + d.timer = time.AfterFunc(wait, func() { + if s.trailing.CompareAndDelete(importID, d) { + s.publish(context.WithoutCancel(ctx), orgID, importID, d.status.Load().(string), true) } - return + }) + if v, pending := s.trailing.LoadOrStore(importID, d); pending { + d.timer.Stop() + v.(*deferredPublish).status.Store(status) } + return } } s.progress.Store(importID, now) s.publisher.PublishMailboxImportProgress(ctx, orgID, importID, status) } +// deferredPublish is a progress event waiting for the end of its one-second window. +type deferredPublish struct { + status atomic.Value + timer *time.Timer +} + // fillVendorPasswords reads a vendor row's passwords from the vendor into the // legs that have none, keeping the servers the row was fixed with. func (s *Service) fillVendorPasswords(ctx context.Context, w repository.ImportWorkRow, p *payload) (string, string) { diff --git a/internal/app/mailboximport/service.go b/internal/app/mailboximport/service.go index b1ab1b65c..215c5c0ab 100644 --- a/internal/app/mailboximport/service.go +++ b/internal/app/mailboximport/service.go @@ -862,17 +862,19 @@ func (s *Service) Dismiss(ctx context.Context, orgID, userID, id uuid.UUID) *err if xerr != nil { return xerr } - if imp.Status == models.ImportRunning { + status := imp.Status + if status == models.ImportRunning { if err := s.repo.Cancel(ctx, orgID, id); err != nil { return errx.InternalError() } - s.audit(ctx, orgID, userID, models.AuditActionUpdate, id, map[string]string{"status": models.ImportCancelled}) + status = models.ImportCancelled + s.audit(ctx, orgID, userID, models.AuditActionUpdate, id, map[string]string{"status": status}) } if err := s.repo.Dismiss(ctx, orgID, id); err != nil { return errx.InternalError() } s.audit(ctx, orgID, userID, models.AuditActionUpdate, id, map[string]string{"dismissed": "true"}) - s.publish(ctx, orgID, id, models.ImportCancelled, true) + s.publish(ctx, orgID, id, status, true) return nil } diff --git a/internal/app/poollink/service.go b/internal/app/poollink/service.go index 6f77b076a..1d2497020 100644 --- a/internal/app/poollink/service.go +++ b/internal/app/poollink/service.go @@ -310,7 +310,11 @@ func (s *service) Plan(ctx context.Context, orgID uuid.UUID) (models.PoolLinkPla if xerr != nil { return models.PoolLinkPlan{}, xerr } - plan := models.PoolLinkPlan{Tier: "free", Enrolled: enrolled, PriceUSD: config.PoolLinkPlanPriceUSD, WarmupEntitled: true} + warming, xerr := s.emails.CountWarmingForOrganization(ctx, orgID) + if xerr != nil { + return models.PoolLinkPlan{}, xerr + } + plan := models.PoolLinkPlan{Tier: "free", Enrolled: enrolled, Warming: warming, PriceUSD: config.PoolLinkPlanPriceUSD, WarmupEntitled: true} if config.BillingProvider() == "none" { plan.Tier = "paid" return plan, nil diff --git a/internal/models/poollink.go b/internal/models/poollink.go index 212e4c2d8..6f6ff9f1e 100644 --- a/internal/models/poollink.go +++ b/internal/models/poollink.go @@ -96,6 +96,8 @@ type PoolLinkPlan struct { // MailboxLimit is nil when unlimited. MailboxLimit *int `json:"mailbox_limit"` Enrolled int `json:"enrolled"` + // Warming is how many of them have warmup on and not paused; connected is not warming. + Warming int `json:"warming"` // PriceUSD is the monthly price of the paid tier, for the upgrade card. PriceUSD int `json:"price_usd"` // UpgradeURL is empty when billing is off or the plan has no price. diff --git a/internal/repository/pg_email.go b/internal/repository/pg_email.go index d1ea1672a..0cb64b37b 100644 --- a/internal/repository/pg_email.go +++ b/internal/repository/pg_email.go @@ -204,6 +204,8 @@ type EmailRepository interface { // CountForOrganization returns the number of email accounts attached to the // given organization. Used by the free-trial inbox cap. CountForOrganization(ctx context.Context, orgID uuid.UUID) (int, *errx.Error) + // CountWarmingForOrganization counts the workspace's active mailboxes with warmup on and not paused. + CountWarmingForOrganization(ctx context.Context, orgID uuid.UUID) (int, *errx.Error) // ListWarmupScheduleCandidates returns active mailboxes that should have a // running warmup chain but currently have no pending warmup task: either @@ -740,6 +742,17 @@ func (r *emailRepository) ListActiveAccountsByWorker(ctx context.Context, worker return ids, rows.Err() } +func (r *emailRepository) CountWarmingForOrganization(ctx context.Context, orgID uuid.UUID) (int, *errx.Error) { + var count int + query := `SELECT COUNT(*) FROM email_accounts + WHERE organization_id = $1 AND status = 'active' AND warmup IS NOT NULL AND warmup_paused_at IS NULL` + if err := r.DB.QueryRow(ctx, query, orgID).Scan(&count); err != nil { + db.CaptureError(err, query, []any{orgID}, "queryrow") + return 0, errx.InternalError() + } + return count, nil +} + func (r *emailRepository) CountForOrganization(ctx context.Context, orgID uuid.UUID) (int, *errx.Error) { var count int query := `SELECT COUNT(*) FROM email_accounts WHERE organization_id = $1` diff --git a/web/src/components/app/emails/CloudPathsPanel.tsx b/web/src/components/app/emails/CloudPathsPanel.tsx index 5cfa33dd2..aa559484b 100644 --- a/web/src/components/app/emails/CloudPathsPanel.tsx +++ b/web/src/components/app/emails/CloudPathsPanel.tsx @@ -23,7 +23,7 @@ export default function CloudPathsPanel({ onAdd, }: { mailboxCount: number; - /** Mailboxes with warmup on and not paused; connected is not the same as warming. */ + /** Mailboxes with warmup on and not paused, from the page's list; used until the server's count loads. */ warmingCount: number; onAdd: () => void; }) { @@ -42,6 +42,9 @@ export default function CloudPathsPanel({ // stands in while it is in flight. const limit = plan ? plan.mailbox_limit : access.paid ? null : FREE_MAILBOXES; const used = plan?.enrolled ?? mailboxCount; + // The server counts the whole workspace; the page's list is filtered and paged, so it is only a stand-in. + const total = plan?.enrolled ?? mailboxCount; + const warming = plan?.warming ?? warmingCount; const free = limit !== null; const allowance = limit ?? FREE_MAILBOXES; const canBuy = free && access.billing && access.isOwner; @@ -99,11 +102,11 @@ export default function CloudPathsPanel({ {free ? `${used} of ${allowance} free mailboxes used. ` : ""} - {warmingCount === 0 + {warming === 0 ? "No mailbox is warming yet." - : `${warmingCount} of ${mailboxCount} mailbox${mailboxCount === 1 ? "" : "es"} warming in the ${onWarmupPlan ? "premium " : ""}pool.`} + : `${warming} of ${total} mailbox${total === 1 ? "" : "es"} warming in the ${onWarmupPlan ? "premium " : ""}pool.`} - {warmingCount === 0 && mailboxCount > 0 && ( + {warming === 0 && total > 0 && ( Turn on warmup from a mailbox's Warmup tab, or select several and start it for all. )} {(linkedLabel || (free && canBuy)) && ( diff --git a/web/src/components/app/emails/import/MailboxImportsMenu.tsx b/web/src/components/app/emails/import/MailboxImportsMenu.tsx index f93ae8687..c04d4d7dc 100644 --- a/web/src/components/app/emails/import/MailboxImportsMenu.tsx +++ b/web/src/components/app/emails/import/MailboxImportsMenu.tsx @@ -47,6 +47,18 @@ function jobState(job: MailboxImport): JobState { tone: "working", }; } + // What needs the person comes first; an authorization still running is mentioned alongside. + const alsoAuthorizing = authorizing > 0 ? ` ${plural(authorizing, "more is", "more are")} still being authorized.` : ""; + if (c.failed > 0) { + return { text: `${c.failed.toLocaleString()} failed`, hint: `Open to see why and retry.${alsoAuthorizing}`, tone: "error" }; + } + if (signin > 0) { + return { + text: `${plural(signin, "mailbox needs", "mailboxes need")} sign-in`, + hint: `Open to sign in to each one.${alsoAuthorizing}`, + tone: "action", + }; + } if (authorizing > 0) { return { text: `Authorizing ${plural(authorizing, "mailbox", "mailboxes")}`, @@ -54,16 +66,6 @@ function jobState(job: MailboxImport): JobState { tone: "waiting", }; } - if (c.failed > 0) { - return { text: `${c.failed.toLocaleString()} failed`, hint: "Open to see why and retry.", tone: "error" }; - } - if (signin > 0) { - return { - text: `${plural(signin, "mailbox needs", "mailboxes need")} sign-in`, - hint: "Open to sign in to each one.", - tone: "action", - }; - } const ok = c.connected + c.updated; return { text: ok > 0 ? `${ok.toLocaleString()} connected` : "Done", tone: "done" }; } diff --git a/web/src/lib/api/models/app/cloudlink/CloudLink.ts b/web/src/lib/api/models/app/cloudlink/CloudLink.ts index c7669885e..23eae631c 100644 --- a/web/src/lib/api/models/app/cloudlink/CloudLink.ts +++ b/web/src/lib/api/models/app/cloudlink/CloudLink.ts @@ -28,6 +28,8 @@ export interface PoolLinkPlan { /** null when unlimited */ mailbox_limit: number | null; enrolled: number; + /** Mailboxes with warmup on and not paused; older servers do not send it. */ + warming?: number; price_usd: number; /** The cloud's billing page with the warmup plan checkout open; only on the free tier. */ upgrade_url?: string;