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

This commit is contained in:
Matthew Meszaros
2026-09-24 18:38:47 +02:00
parent d8e99877b2
commit 561ba07b35
8 changed files with 72 additions and 29 deletions
+26 -11
View File
@@ -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) {
+5 -3
View File
@@ -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
}
+5 -1
View File
@@ -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
+2
View File
@@ -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.
+13
View File
@@ -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`
@@ -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({
<span className="min-w-0 flex-1 leading-snug">
<span className="font-medium">
{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.`}
</span>
{warmingCount === 0 && mailboxCount > 0 && (
{warming === 0 && total > 0 && (
<span className="text-slate-500"> Turn on warmup from a mailbox&apos;s Warmup tab, or select several and start it for all.</span>
)}
{(linkedLabel || (free && canBuy)) && (
@@ -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" };
}
@@ -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;