From 561ba07b35153a8648528bb213cd5cd4a5adaa7a Mon Sep 17 00:00:00 2001 From: Matthew Meszaros Date: Thu, 24 Sep 2026 18:38:47 +0200 Subject: [PATCH] feat: send a deferred import progress event with the latest status and cancel it when a final one goes out, report a hidden import's real status instead of cancelled, rank failed and sign-in rows above authorizing ones in the imports menu and its badge, and count warming and connected mailboxes for the pool banner on the server (warming on the pool plan) rather than from the filtered, paged list --- internal/app/mailboximport/runner.go | 37 +++++++++++++------ internal/app/mailboximport/service.go | 8 ++-- internal/app/poollink/service.go | 6 ++- internal/models/poollink.go | 2 + internal/repository/pg_email.go | 13 +++++++ .../components/app/emails/CloudPathsPanel.tsx | 11 ++++-- .../app/emails/import/MailboxImportsMenu.tsx | 22 ++++++----- .../lib/api/models/app/cloudlink/CloudLink.ts | 2 + 8 files changed, 72 insertions(+), 29 deletions(-) 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;