diff --git a/cmd/backend/main.go b/cmd/backend/main.go index fe67fbf33..8cc56a0cc 100644 --- a/cmd/backend/main.go +++ b/cmd/backend/main.go @@ -2292,7 +2292,7 @@ func main() { drained.Add(1) go func() { defer drained.Done() - if !mailboxImportService.Drain(shutdownGrace) { + if !mailboxImportService.Drain(shutdownCtx) { log.Println("Import rows still connecting at shutdown; they resume on the next instance") } }() diff --git a/docs/content/docs/development/install.mdx b/docs/content/docs/development/install.mdx index 1609aea30..5bdf0f6b9 100644 --- a/docs/content/docs/development/install.mdx +++ b/docs/content/docs/development/install.mdx @@ -300,7 +300,7 @@ By hand, from the install directory: docker compose -p warmbly pull && docker compose -p warmbly up -d ``` -An update does not cut off work in progress. The generated stack gives each application container 60 seconds to stop (`stop_grace_period`), and in that time the backend finishes a mailbox connect already under way and the import rows it is connecting, and a worker finishes the send it is on. Anything still unfinished after that is picked up again when the new containers start. See [updates without lost work](/development/deployment-guide/#updates-without-lost-work). +An update does not cut off work in progress. The generated stack gives each application container 60 seconds to stop (`stop_grace_period`), and in that time the backend finishes a mailbox connect already under way and the import rows it is connecting, and a worker finishes the send it is on. Anything still unfinished after that is picked up again by the new containers: a mailbox a worker was loading on its first heartbeat, and an import row once its 2-minute lease runs out. See [updates without lost work](/development/deployment-guide/#updates-without-lost-work). ## Removing it diff --git a/internal/app/mailboximport/retry_test.go b/internal/app/mailboximport/retry_test.go index fd2e0d4bf..adff23358 100644 --- a/internal/app/mailboximport/retry_test.go +++ b/internal/app/mailboximport/retry_test.go @@ -93,19 +93,20 @@ func TestGrantForHostMatchesTheRowsProviderAndDomain(t *testing.T) { } } -func TestDrainWaitsForRowsInFlight(t *testing.T) { - s := &Service{} - s.inflight.Add(1) +func TestDrainWaitsForTheRunnerToStop(t *testing.T) { + s := &Service{stopped: make(chan struct{})} + short, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond) + defer cancel() + if s.Drain(short) { + t.Fatal("drain reported a runner still working as stopped") + } go func() { time.Sleep(50 * time.Millisecond) - s.inflight.Done() + close(s.stopped) }() - if !s.Drain(2 * time.Second) { - t.Fatal("drain gave up on a row that finished") - } - s.inflight.Add(1) - defer s.inflight.Done() - if s.Drain(50 * time.Millisecond) { - t.Fatal("drain reported a row still connecting as finished") + ctx, cancel2 := context.WithTimeout(context.Background(), 2*time.Second) + defer cancel2() + if !s.Drain(ctx) { + t.Fatal("drain gave up on a runner that stopped") } } diff --git a/internal/app/mailboximport/runner.go b/internal/app/mailboximport/runner.go index a7ccb7129..2bf656350 100644 --- a/internal/app/mailboximport/runner.go +++ b/internal/app/mailboximport/runner.go @@ -36,8 +36,10 @@ func (s *Service) Kick() { } } -// Start runs the import loop until ctx ends, plus the slower upkeep loops. +// Start runs the import loop until ctx ends, plus the slower upkeep loops. It +// returns only once the pass in progress has finished or handed back its rows. func (s *Service) Start(ctx context.Context) { + defer close(s.stopped) go jobrun.Loop(ctx, "mailbox_import_upkeep", time.Minute, true, func(ctx context.Context) error { s.reconcileSignins(ctx) s.resumeVendorAuthorizations(ctx) @@ -92,9 +94,7 @@ func (s *Service) pass(ctx context.Context) { sem := make(chan struct{}, config.MailboxImportConcurrency) for _, w := range rows { wg.Add(1) - s.inflight.Add(1) go func(w repository.ImportWorkRow) { - defer s.inflight.Done() defer wg.Done() defer func() { if rec := recover(); rec != nil { @@ -124,17 +124,14 @@ func (s *Service) pass(ctx context.Context) { s.completeFinished(work) } -// Drain waits for the rows being connected to finish, up to timeout; false when some did not. -func (s *Service) Drain(timeout time.Duration) bool { - done := make(chan struct{}) - go func() { - s.inflight.Wait() - close(done) - }() +// Drain waits until the runner has stopped (ctx for Start ended and the pass in +// progress finished its rows, or handed back the ones it never started), or +// until ctx ends; false when ctx ended first. +func (s *Service) Drain(ctx context.Context) bool { select { - case <-done: + case <-s.stopped: return true - case <-time.After(timeout): + case <-ctx.Done(): return false } } @@ -452,29 +449,12 @@ func (s *Service) grantForHost(ctx context.Context, orgID uuid.UUID, mailHost, e // resumeGrantedSignins queues again the rows waiting on sign-in whose domain an // administrator's grant now covers, so one admin approval finishes all of them. func (s *Service) resumeGrantedSignins(ctx context.Context) { - if s.delegator == nil { - return - } - rows, err := s.repo.SigninRows(ctx, 1000) + rows, err := s.repo.CoveredSigninRows(ctx, 1000) if err != nil || len(rows) == 0 { return } - type key struct { - org uuid.UUID - host, domain string - } - covered := map[key]bool{} resumed := false for _, w := range rows { - k := key{w.OrgID, w.MailHost, domainOf(w.Email)} - ok, seen := covered[k] - if !seen { - ok = s.grantForHost(ctx, w.OrgID, w.MailHost, w.Email) != nil - covered[k] = ok - } - if !ok { - continue - } if err := s.repo.ResumeParked(ctx, w.ImportID, w.Line, w.Code); err == nil { resumed = true } diff --git a/internal/app/mailboximport/service.go b/internal/app/mailboximport/service.go index 215c5c0ab..a975ce20e 100644 --- a/internal/app/mailboximport/service.go +++ b/internal/app/mailboximport/service.go @@ -146,10 +146,10 @@ type Service struct { redirected sync.Map kick chan struct{} - inflight sync.WaitGroup // rows being connected, so a shutdown can let them finish - progress sync.Map // import id -> time.Time of the last progress event - trailing sync.Map // import id -> a progress event deferred to the end of its window - causes sync.Map // scrubbed server reply -> refined cause key + stopped chan struct{} // closed when Start returns, after the pass in progress has finished its rows + progress sync.Map // import id -> time.Time of the last progress event + trailing sync.Map // import id -> a progress event deferred to the end of its window + causes sync.Map // scrubbed server reply -> refined cause key } func NewService(d Deps) *Service { @@ -164,7 +164,7 @@ func NewService(d Deps) *Service { detector: d.Detector, allowance: d.Allowance, asker: d.Asker, warmup: d.Warmup, auditor: d.Auditor, publisher: d.Publisher, google: d.GoogleSignin, delegator: d.Delegator, vendors: d.Vendors, domains: d.Domains, - kick: make(chan struct{}, 1), + kick: make(chan struct{}, 1), stopped: make(chan struct{}), } } diff --git a/internal/infrastructure/eventbus/nats.go b/internal/infrastructure/eventbus/nats.go index 13ec74b07..01dbf61b7 100644 --- a/internal/infrastructure/eventbus/nats.go +++ b/internal/infrastructure/eventbus/nats.go @@ -303,6 +303,7 @@ func (b *NATSBus) Subscribe(ctx context.Context, topics []string, group string, if b.closed { b.mu.Unlock() cc.Stop() + <-cc.Closed() return ErrBusClosed } b.subscribers = append(b.subscribers, cc) @@ -311,6 +312,8 @@ func (b *NATSBus) Subscribe(ctx context.Context, topics []string, group string, // Block until ctx is cancelled, mirroring KafkaBus.Subscribe semantics. <-ctx.Done() cc.Stop() + // Closed fires once the callback in progress has returned, so shutdown waits for it. + <-cc.Closed() return ctx.Err() } @@ -327,6 +330,7 @@ func (b *NATSBus) Close() error { for _, s := range subs { s.Stop() + <-s.Closed() } if b.nc != nil { // Drain rather than hard-close to flush in-flight publishes. diff --git a/internal/models/mailbox_import.go b/internal/models/mailbox_import.go index 3d17335d1..3f6177651 100644 --- a/internal/models/mailbox_import.go +++ b/internal/models/mailbox_import.go @@ -228,6 +228,8 @@ type MailboxImport struct { Total int `json:"total"` Counts MailboxImportCounts `json:"counts"` Causes []MailboxImportCause `json:"causes"` + // Authorizing counts the rows an inbox vendor is still authorizing, by provider ("google", "microsoft"). + Authorizing map[string]int `json:"authorizing"` } type MailboxImportRow struct { diff --git a/internal/repository/pg_mailbox_import.go b/internal/repository/pg_mailbox_import.go index c266f41a1..0a414a3f2 100644 --- a/internal/repository/pg_mailbox_import.go +++ b/internal/repository/pg_mailbox_import.go @@ -113,9 +113,9 @@ type MailboxImportRepository interface { Dismiss(ctx context.Context, orgID, id uuid.UUID) error // TouchParked moves rows still waiting to the back of ParkedRows, so no import holds the others back. TouchParked(ctx context.Context, cause string, rows []ImportWorkRow) error - // SigninRows lists rows waiting on a Google or Microsoft sign-in that still hold their settings, in imports not cancelled. - // Code carries the row's cause. - SigninRows(ctx context.Context, limit int) ([]ImportWorkRow, error) + // CoveredSigninRows lists rows waiting on a Google or Microsoft sign-in, in imports not cancelled, + // that an active grant of their workspace now covers. Code carries the row's cause. + CoveredSigninRows(ctx context.Context, limit int) ([]ImportWorkRow, error) // ResumeParked queues a parked row again and reopens its import when it had completed. ResumeParked(ctx context.Context, id uuid.UUID, line int, cause string) error GetMapping(ctx context.Context, orgID uuid.UUID, signature string) (models.MailboxImportMapping, bool, error) @@ -234,6 +234,32 @@ func (r *mailboxImportRepository) fillCounts(ctx context.Context, imp *models.Ma return err } + query = ` + SELECT CASE code WHEN 'google_signin' THEN 'google' ELSE 'microsoft' END, count(*) + FROM mailbox_import_rows + WHERE import_id = $1 AND status = 'needs_signin' AND cause = '` + ParkedVendorCause + `' + GROUP BY 1` + arows, err := r.DB.Query(ctx, query, imp.ID) + if err != nil { + db.CaptureError(err, query, nil, "query") + return err + } + imp.Authorizing = map[string]int{} + for arows.Next() { + var provider string + var n int + if err := arows.Scan(&provider, &n); err != nil { + arows.Close() + db.CaptureError(err, query, nil, "scan") + return err + } + imp.Authorizing[provider] = n + } + arows.Close() + if err := arows.Err(); err != nil { + return err + } + query = ` SELECT cause, count(*), bool_or(payload <> '' AND status = 'failed') FROM mailbox_import_rows @@ -733,13 +759,20 @@ func (r *mailboxImportRepository) TouchParked(ctx context.Context, cause string, return nil } -func (r *mailboxImportRepository) SigninRows(ctx context.Context, limit int) ([]ImportWorkRow, error) { +func (r *mailboxImportRepository) CoveredSigninRows(ctx context.Context, limit int) ([]ImportWorkRow, error) { + // Coverage is decided here, so rows no grant covers can never crowd out the ones it does. query := ` SELECT r.import_id, i.organization_id, r.line, r.email, r.mail_host, r.cause FROM mailbox_import_rows r JOIN mailbox_imports i ON i.id = r.import_id WHERE r.status = 'needs_signin' AND r.cause IN ('microsoft_signin', 'google_signin') AND r.payload <> '' AND i.status <> 'cancelled' + AND EXISTS ( + SELECT 1 FROM mailbox_domain_grants g + WHERE g.organization_id = i.organization_id AND g.status = 'active' + AND g.provider = CASE r.cause WHEN 'microsoft_signin' THEN 'microsoft' ELSE 'google' END + AND lower(split_part(r.email, '@', 2)) = ANY(g.domains) + ) ORDER BY r.updated_at LIMIT $1` rows, err := r.DB.Query(ctx, query, limit) diff --git a/web/src/app/app/emails/page.tsx b/web/src/app/app/emails/page.tsx index 1b0b50d2d..cc3ccb907 100644 --- a/web/src/app/app/emails/page.tsx +++ b/web/src/app/app/emails/page.tsx @@ -368,7 +368,7 @@ export default function AddressesPage() { openMigration()} /> setCloudDialog(true)} mailboxCount={stats.total} /> - {!emailsData.isLoading && p?.setAddEmail(true)} />} + {!emailsData.isLoading && p?.setAddEmail(true)} />} {/* Hosted, the pool is thousands of mailboxes: the pool-size advice is self-host only. */} {cloud.selfHosted && ( void; -}) { +export default function CloudPathsPanel({ mailboxCount, onAdd }: { mailboxCount: number; onAdd: () => void }) { const authConfig = useAuthConfig(); const access = useFeatureAccess(); const hosted = authConfig.data?.self_hosted === false; @@ -42,9 +33,10 @@ 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. + // The server counts the whole workspace; the page's list is filtered and paged, so it only stands in for the total. const total = plan?.enrolled ?? mailboxCount; - const warming = plan?.warming ?? warmingCount; + // Only the server knows how many warm across the workspace; without it the banner says nothing about warming. + const warming = plan?.warming; const free = limit !== null; const allowance = limit ?? FREE_MAILBOXES; const canBuy = free && access.billing && access.isOwner; @@ -53,7 +45,7 @@ export default function CloudPathsPanel({ const dialog = setPlanOpen(false)} />; - if (mailboxCount === 0 && linked === 0) { + if (total === 0 && linked === 0) { return (
@@ -102,9 +94,13 @@ export default function CloudPathsPanel({ {free ? `${used} of ${allowance} free mailboxes used. ` : ""} - {warming === 0 - ? "No mailbox is warming yet." - : `${warming} of ${total} mailbox${total === 1 ? "" : "es"} warming in the ${onWarmupPlan ? "premium " : ""}pool.`} + {warming === undefined + ? free + ? "" + : `${total} mailbox${total === 1 ? "" : "es"} connected.` + : warming === 0 + ? "No mailbox is warming yet." + : `${warming} of ${total} mailbox${total === 1 ? "" : "es"} warming in the ${onWarmupPlan ? "premium " : ""}pool.`} {warming === 0 && total > 0 && ( Turn on warmup from a mailbox's Warmup tab, or select several and start it for all. diff --git a/web/src/components/app/emails/import/MailboxImportsMenu.tsx b/web/src/components/app/emails/import/MailboxImportsMenu.tsx index c04d4d7dc..2de1c7528 100644 --- a/web/src/components/app/emails/import/MailboxImportsMenu.tsx +++ b/web/src/components/app/emails/import/MailboxImportsMenu.tsx @@ -40,25 +40,28 @@ function jobState(job: MailboxImport): JobState { const settled = job.total - inFlight; if (job.status === "cancelled") return { text: "Stopped", tone: "muted" }; - if (job.status === "running" && inFlight > 0) { - return { - text: `Connecting ${settled.toLocaleString()} of ${job.total.toLocaleString()}`, - hint: "Each mailbox is checked against its mail server, a few seconds each.", - 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.` : ""; + // What needs the person comes first; work still going on is mentioned alongside. + const connecting = job.status === "running" && inFlight > 0; + const alongside = + (connecting ? ` Still connecting ${settled.toLocaleString()} of ${job.total.toLocaleString()}.` : "") + + (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" }; + return { text: `${c.failed.toLocaleString()} failed`, hint: `Open to see why and retry.${alongside}`, tone: "error" }; } if (signin > 0) { return { text: `${plural(signin, "mailbox needs", "mailboxes need")} sign-in`, - hint: `Open to sign in to each one.${alsoAuthorizing}`, + hint: `Open to sign in to each one.${alongside}`, tone: "action", }; } + if (connecting) { + return { + text: `Connecting ${settled.toLocaleString()} of ${job.total.toLocaleString()}`, + hint: "Each mailbox is checked against its mail server, a few seconds each.", + tone: "working", + }; + } if (authorizing > 0) { return { text: `Authorizing ${plural(authorizing, "mailbox", "mailboxes")}`, diff --git a/web/src/components/app/emails/import/RunStep.tsx b/web/src/components/app/emails/import/RunStep.tsx index 2ff6ecb85..97b0c646d 100644 --- a/web/src/components/app/emails/import/RunStep.tsx +++ b/web/src/components/app/emails/import/RunStep.tsx @@ -226,11 +226,7 @@ export default function RunStep({ const signinWaiting = Math.max(0, c.needs_signin - authorizing); const failureCauses = data.causes.filter((x) => x.cause !== VENDOR_AUTHORIZING); // Microsoft rows the vendor is authorizing can be finished sooner by one admin approval. - const msAuthorizing = - authorizing > 0 && - (rows.data?.data ?? []).some( - (r) => r.cause === VENDOR_AUTHORIZING && (r.code === "microsoft_signin" || mailHostOAuthProvider(r.mail_host) === "outlook"), - ); + const msAuthorizing = data.authorizing?.microsoft ?? 0; const allowanceHit = data.causes.some((x) => x.cause === "allowance_reached"); const msSignin = msGrants && c.needs_signin > 0 && data.causes.some((x) => x.cause === "microsoft_signin"); // Rows parked on the vendor keep the import running; only they left means it is waiting, not importing. @@ -368,7 +364,7 @@ export default function RunStep({

Want it faster?

- {msAuthorizing && msGrants && !granted && ( + {msAuthorizing > 0 && msGrants && !granted && (

One sign-in with the domain's Global Administrator account ({vendorLabel(data.vendor) || "your inbox vendor"}{" "} @@ -389,7 +385,7 @@ export default function RunStep({

)}

- {msAuthorizing && msGrants && !granted ? "Or connect them one by one:" : "Connect them one by one:"} Sign in on + {msAuthorizing > 0 && msGrants && !granted ? "Or connect them one by one:" : "Connect them one by one:"} Sign in on a row signs in as that mailbox and connects only it.

diff --git a/web/src/lib/api/models/app/emails/MailboxImport.ts b/web/src/lib/api/models/app/emails/MailboxImport.ts index 3a941fca2..d85ef87ad 100644 --- a/web/src/lib/api/models/app/emails/MailboxImport.ts +++ b/web/src/lib/api/models/app/emails/MailboxImport.ts @@ -239,6 +239,8 @@ export interface MailboxImport { total: number; counts: MailboxImportCounts; causes: ImportCause[]; + /** Rows an inbox vendor is still authorizing, by provider; older servers do not send it. */ + authorizing?: Partial>; } export interface ImportRow {