feat: make the backend drain wait for the import runner itself to stop on the shutdown deadline, resume only sign-in rows an active grant covers (decided in SQL so uncovered rows cannot crowd them out), wait for NATS consumers to close after stopping them, count rows a vendor is authorizing by provider on the import so the admin sign-in button counts and shows only Microsoft rows, keep failed and sign-in imports ahead of connecting ones in the imports menu, and base the pool banner's empty state and warming count on workspace totals only

This commit is contained in:
Matthew Meszaros
2026-09-24 18:57:01 +02:00
parent 561ba07b35
commit f0861fa129
13 changed files with 105 additions and 88 deletions
+1 -1
View File
@@ -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")
}
}()
+1 -1
View File
@@ -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
+12 -11
View File
@@ -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")
}
}
+10 -30
View File
@@ -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
}
+5 -5
View File
@@ -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{}),
}
}
+4
View File
@@ -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.
+2
View File
@@ -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 {
+37 -4
View File
@@ -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)
+1 -1
View File
@@ -368,7 +368,7 @@ export default function AddressesPage() {
<AdvisorSummaryBar surface="emails" noun="mailbox" nounPlural="mailboxes" />
<SigninMigrationBanner total={migration.data?.total ?? 0} onOpen={() => openMigration()} />
<CloudPoolBanner onConnect={() => setCloudDialog(true)} mailboxCount={stats.total} />
{!emailsData.isLoading && <CloudPathsPanel mailboxCount={stats.total} warmingCount={warmupActive} onAdd={() => p?.setAddEmail(true)} />}
{!emailsData.isLoading && <CloudPathsPanel mailboxCount={stats.total} onAdd={() => p?.setAddEmail(true)} />}
{/* Hosted, the pool is thousands of mailboxes: the pool-size advice is self-host only. */}
{cloud.selfHosted && (
<WarmupCoverageNotice
@@ -17,16 +17,7 @@ import WarmupPlanDialog from "@/components/app/billing/WarmupPlanDialog";
const FREE_MAILBOXES = 10;
export default function CloudPathsPanel({
mailboxCount,
warmingCount,
onAdd,
}: {
mailboxCount: number;
/** Mailboxes with warmup on and not paused, from the page's list; used until the server's count loads. */
warmingCount: number;
onAdd: () => 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 = <WarmupPlanDialog open={planOpen} onClose={() => setPlanOpen(false)} />;
if (mailboxCount === 0 && linked === 0) {
if (total === 0 && linked === 0) {
return (
<div className="py-3">
<motion.div initial={{ opacity: 0, y: 8 }} animate={{ opacity: 1, y: 0 }} className="rounded-xl border border-slate-200 overflow-hidden">
@@ -102,9 +94,13 @@ export default function CloudPathsPanel({
<span className="min-w-0 flex-1 leading-snug">
<span className="font-medium">
{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.`}
</span>
{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>
@@ -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")}`,
@@ -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({
</p>
<div className="mt-2 pt-2 border-t border-sky-200/70">
<p className="text-[11.5px] font-medium text-sky-900">Want it faster?</p>
{msAuthorizing && msGrants && !granted && (
{msAuthorizing > 0 && msGrants && !granted && (
<div className="mt-1.5">
<button
type="button"
@@ -379,7 +375,7 @@ export default function RunStep({
{consent.busy ? <Loader2Icon className="w-3 h-3 animate-spin" /> : <Building2Icon className="w-3 h-3" />}
{consent.busy
? "Waiting for the administrator…"
: `Connect all ${authorizing.toLocaleString()} at once (admin sign-in)`}
: `Connect all ${msAuthorizing.toLocaleString()} Microsoft mailbox${msAuthorizing === 1 ? "" : "es"} at once (admin sign-in)`}
</button>
<p className="text-[11px] text-sky-800/90 leading-relaxed mt-1">
One sign-in with the domain&apos;s Global Administrator account ({vendorLabel(data.vendor) || "your inbox vendor"}{" "}
@@ -389,7 +385,7 @@ export default function RunStep({
</div>
)}
<p className="text-[11px] text-sky-800/90 leading-relaxed mt-1">
{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.
</p>
</div>
@@ -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<Record<"google" | "microsoft", number>>;
}
export interface ImportRow {