mirror of
https://github.com/warmbly/warmbly.git
synced 2026-10-03 08:02:04 +00:00
feat: keep mailbox credential checks and imports from timing out behind a worker's command queue by loading mailboxes off the bus loop (a republish never dials twice, commands wait for an in-flight load, a failed load raises the account's auth or server error), storing Kafka offsets for background commit instead of a synchronous commit per message, reconciling moved, unplaced and dead-worker mailboxes at once while spreading the safety-net republish over 30 minutes, sending checks over each worker's Redis channel with failover in placement order, retrying unanswered import rows within the lease, finishing a connect the browser left and an import row a restart interrupted, and connecting InboxKit Google and Microsoft mailboxes with no sign-in through a vendor-authorized grant that covers only the domains the vendor account lists
This commit is contained in:
@@ -1692,6 +1692,7 @@ func main() {
|
||||
Cipher: cipherService,
|
||||
Mailboxes: emailRepostory,
|
||||
Reconnect: emailService,
|
||||
Grants: delegationService,
|
||||
Cache: cache,
|
||||
NewClient: sandboxVendors,
|
||||
})
|
||||
|
||||
@@ -198,6 +198,7 @@ func main() {
|
||||
// the rolling 1m counters into a WorkerHealth event, publishes via the
|
||||
// event bus so the consumer can write a row into worker_health_samples.
|
||||
go workerService.Heartbeat(ctx)
|
||||
go workerService.ListenValidations(ctx)
|
||||
go workerService.RunHealth(ctx, 30*time.Second)
|
||||
|
||||
// heartbeatDone closes once the farewell beat has been sent. Main waits on
|
||||
|
||||
@@ -591,7 +591,7 @@ A `400` whose `code` is `mailbox_validation_timeout` comes from the SMTP and IMA
|
||||
}
|
||||
```
|
||||
|
||||
The same code with a message that begins `Warmbly's worker did not report back` means no worker answered the check at all, so the mail server was never tested. On the hosted product, retry and contact support if it persists. On a self-hosted instance, check that a worker is running, heartbeating, and reaching the same Redis as the backend.
|
||||
The same code with a message that begins `Warmbly's worker did not report back` means no worker answered the check at all, so the mail server was never tested. The check is offered to each healthy worker in turn and waits on at most two that take it without answering, so this means the fleet as a whole is not answering. On the hosted product, retry and contact support if it persists. On a self-hosted instance, check that a worker is running, heartbeating, and reaching the same Redis as the backend: the check reaches the worker over Redis and its answer comes back the same way.
|
||||
|
||||
**How to fix:**
|
||||
- Read which leg stayed silent. One leg passing and the other hanging is the network between the worker and that port, not the password: a port that nothing listens on, or that a firewall drops, takes the full timeout rather than refusing straight away
|
||||
|
||||
@@ -46,7 +46,7 @@ flowchart LR
|
||||
|
||||
## Data flow
|
||||
|
||||
The frontend talks to the backend over REST + JWT, and to the realtime service over WebSocket. Backend writes business state to Postgres. Backend, tracking, and workers publish to the event bus: NATS JetStream with JSON encoding by default, Kafka with Avro and Schema Registry as an opt-in build. The consumer reads the bus and updates Postgres (analytics, suppression, deliverability). Workers subscribe to a topic named for their worker UUID (`w.<uuid>`) and publish results to `jobs.worker-events`; see the [event system](/development/events/) for the full topic map.
|
||||
The frontend talks to the backend over REST + JWT, and to the realtime service over WebSocket. Backend writes business state to Postgres. Backend, tracking, and workers publish to the event bus: NATS JetStream with JSON encoding by default, Kafka with Avro and Schema Registry as an opt-in build. The consumer reads the bus and updates Postgres (analytics, suppression, deliverability). Workers subscribe to a topic named for their worker UUID (`w.<uuid>`) and publish results to `jobs.worker-events`; see the [event system](/development/events/) for the full topic map. The one exception is the credential check run before a mailbox is saved: someone is waiting on it, and a worker reads its topic one command at a time, so the backend offers it on the worker's Redis channel `email_validation_request:<uuid>` instead and gets the answer back on Redis. A worker that is not subscribed there (an older build) still gets the check over its topic.
|
||||
|
||||
Realtime fanout is a separate channel: backend and consumer publish JSON events over Redis pub/sub (or Google Cloud Pub/Sub when `PUBSUB_ENABLED=true`); the Elixir realtime service subscribes and pushes to connected WebSocket clients.
|
||||
|
||||
|
||||
@@ -792,7 +792,7 @@ These are the only settings a browser can change, and no environment variable ow
|
||||
| `deliverability.enforce_domain_auth` | boolean | `true` | Whether a sending domain that fails SPF or DMARC stops cold campaign sending and warmup sending from every mailbox on it. Off keeps the check running and still shows the state and the Advisor card, it just never blocks |
|
||||
| `deliverability.auth_grace_hours` | integer, 1 to 720 | `72` | How long a domain must stay failing before the gate applies. The clock starts when the background check first sees the failure, so this is also how much warning the owner gets |
|
||||
|
||||
The four `sync.*` values are read by the backend when a mailbox is loaded onto a worker (on connect, on reassignment, and by the reconciler's periodic republish), so a change reaches every mailbox within a few minutes without a restart. The fixed pacing numbers around them (burst per five minutes, hourly, backfill pace, the flood threshold and the chronic-overage rule) are compiled constants listed under **Instance > Configuration > Effective limits**; see [Mailboxes](/guides/mailboxes/#what-gets-synced) for how the budgets behave. The four `sync.*` values also have a companion read view on **Operations > Sync**, which shows each mailbox's backfill progress and fair-use throttle against them and can clear a throttle or restart a backfill.
|
||||
The four `sync.*` values are read by the backend when a mailbox is loaded onto a worker (on connect, on reassignment, and by the reconciler's periodic republish), so a change reaches every mailbox within 30 minutes without a restart (the republish gives each mailbox one turn per 30 minutes, spread out so no worker gets them all at once). The fixed pacing numbers around them (burst per five minutes, hourly, backfill pace, the flood threshold and the chronic-overage rule) are compiled constants listed under **Instance > Configuration > Effective limits**; see [Mailboxes](/guides/mailboxes/#what-gets-synced) for how the budgets behave. The four `sync.*` values also have a companion read view on **Operations > Sync**, which shows each mailbox's backfill progress and fair-use throttle against them and can clear a throttle or restart a backfill.
|
||||
|
||||
The five `retention.*` values are read by the pruning sweeps on every pass, so shortening one takes effect on the next sweep rather than at the next restart. Deletion is permanent and there is no grace period: what already sits outside a shortened window goes on that sweep. The warmup mail window is the one that reaches into customers' mailboxes: the consumer retires each message and the worker holding the mailbox deletes it there, so shortening it empties the warmup folders of the whole instance down to the new window within a few hours. See [data control](/development/data-control/#what-is-kept-and-for-how-long) for what each log holds and what a shorter window costs.
|
||||
|
||||
|
||||
@@ -56,6 +56,12 @@ The backend drives workers with a `{type, body}` envelope (`internal/models/even
|
||||
|
||||
Types: `SEND_EMAIL`, `ADD_EMAIL`, `REMOVE_EMAIL`, `EMAIL_VALIDATION`, `WARMUP_ACTION`, `MESSAGE_SEEN`, `MAILBOX_IDENTITY`. Each has a typed body struct in `internal/models/` (for example `models.SendEmail` for `SEND_EMAIL`).
|
||||
|
||||
A worker reads its topic one command at a time, so nothing slow runs on that loop. `ADD_EMAIL` connects the mailbox in the background (a republish for a mailbox already loading is ignored, and a command for it waits for that load), and a mailbox that cannot be loaded raises the same `EMAIL_AUTH_ERROR` or `EMAIL_SERVER_ERROR` a failed sync would, so it shows the error instead of looking connected. On Kafka, a handled message's offset is stored and committed in the background rather than with a synchronous commit per message.
|
||||
|
||||
The backend's worker reconciler ships `ADD_EMAIL` at once for a mailbox with no worker, on a worker that stopped heartbeating, or moved to another worker. Every other active mailbox is re-shipped once every 30 minutes as a safety net, each at its own point in the interval.
|
||||
|
||||
`EMAIL_VALIDATION` is normally not sent on the topic at all. The backend publishes the same body to the worker's Redis channel `email_validation_request:<worker uuid>`, which a worker reads alongside its topic, so a check never waits behind queued sends. Redis reports how many subscribers took it; with none, the next healthy worker is tried, and only when no worker listens there does the check go on the topic.
|
||||
|
||||
`MESSAGE_SEEN` carries a read or unread change made in the unified inbox out to the mailbox provider, one event per mailbox and at most `models.SeenRelayChunk` messages each. It answers with nothing: the store was written before the event was published, so the provider's copy is the only thing it changes, and a failure is logged rather than retried because the next sync reports whatever the provider actually holds. Each message carries all three providers' handles and the worker takes the ones its client uses: a provider message id for Gmail and Graph, an IMAP folder plus UID, and the immutable RFC Message-ID, which Graph re-resolves the live id from because a Graph id changes whenever a message moves. The state relayed is the one the row holds when the relay reads it back, not the one the request asked for, so a conflicting toggle leaves the provider agreeing with the store rather than with whoever published last.
|
||||
|
||||
`MAILBOX_IDENTITY` asks the worker holding a mailbox to read its send-as addresses, and one signature, from the provider. It answers on a Redis channel named by the request's `process_id` rather than on the worker-events topic, the same round trip `EMAIL_VALIDATION` uses, because the caller is an HTTP request waiting for it. The control plane makes this call itself in exactly one place, the OAuth handshake, where the token came from the consent the customer just completed and the mailbox has no worker yet.
|
||||
|
||||
@@ -123,7 +123,8 @@ A key that reaches several workspaces (or, at ScaledMail, organizations) lists t
|
||||
|
||||
One vendor account lists up to 10,000 mailboxes.
|
||||
|
||||
- **Credentials are read when each row is worked.** Picking mailboxes stores only which ones you picked. The password, app password and servers of each are fetched from the vendor at the moment its row connects, and the row is then judged exactly like a file row: Google mailboxes need the app password the vendor holds (or [an admin grant](#connect-a-whole-google-workspace-domain) over the domain), and Microsoft mailboxes wait for sign-in unless [a grant](#connect-a-whole-microsoft-365-organization) covers them.
|
||||
- **Credentials are read when each row is worked.** Picking mailboxes stores only which ones you picked. The password, app password and servers of each are fetched from the vendor at the moment its row connects, and the row is then judged exactly like a file row: Google mailboxes need the app password the vendor holds (or [an admin grant](#connect-a-whole-google-workspace-domain) over the domain), and Microsoft mailboxes connect through [a grant](#connect-a-whole-microsoft-365-organization) over their domain.
|
||||
- **The vendor can authorize Warmbly for you.** When a row would otherwise wait for a sign-in and the vendor can approve apps through the admin mailbox it holds on the domain (InboxKit can), Warmbly asks it to: tenant-wide consent for Warmbly's Microsoft app, or domain-wide delegation for Warmbly's Google client. The rows show as **Authorizing** in the meantime, which usually takes a few minutes, and the import stays open until they connect on their own, with nobody signing in. One request covers every mailbox on the domain. The resulting grant covers only domains this vendor account holds, even when the vendor's Microsoft organization has others. A Google domain needs an admin mailbox listed at the vendor. When the vendor cannot authorize the domain, or has not finished within two hours, its rows fall back to **Sign in** with the reason on the row. Cancelling the import does the same for rows still waiting.
|
||||
- **The key is sealed and never shown again.** The API key is sealed with the workspace's own encryption key. The list of connections shows each one's vendor, label, status and mailbox count, never the key. **Update** replaces the key after checking the new one with the vendor.
|
||||
- **Passwords the vendor rotates are picked up.** An SMTP and IMAP mailbox that came from a vendor and has stopped with an unresolved connection error is looked at every 15 minutes. When the vendor now holds a different password for it, Warmbly verifies that password against the server and reconnects the mailbox with it. Each mailbox is tried at most once every 6 hours, so a password that still fails is not retried on every pass.
|
||||
- **A refused key marks the connection.** A key the vendor accepts but that reaches no workspace is refused with `mailbox_vendor_no_workspace`, since there is nothing to list. When the vendor stops accepting the key, the connection shows as invalid with the vendor's reason, and rows that needed it fail as `vendor_unauthorized` until the key is updated.
|
||||
@@ -245,9 +246,9 @@ A setting that cannot be applied does not fail the row: the mailbox connects and
|
||||
|
||||
## While it runs
|
||||
|
||||
The import runs on the server. Closing the tab does not stop it, and everyone in the workspace who can manage mailboxes sees its progress live, with the mailbox list filling in behind it. **Recent imports** in the dialog reopens one later.
|
||||
The import runs on the server. Closing the tab does not stop it, and everyone in the workspace who can manage mailboxes sees its progress live, with the mailbox list filling in behind it. A server restart does not lose rows either: a row being connected at the time is picked up again, and a mailbox it had already created is finished with the import's settings rather than skipped as one that was already there. **Recent imports** in the dialog reopens one later.
|
||||
|
||||
Rows connect 8 at a time, and at most 3 at once against any one mail host, since providers throttle bursts of sign-ins from one address. Each credential is verified against the live server before it is saved, exactly like a single connect, so wrong credentials fail here rather than at the first send.
|
||||
Rows connect 8 at a time, and at most 3 at once against any one mail host, since providers throttle bursts of sign-ins from one address. Each credential is verified against the live server before it is saved, exactly like a single connect, so wrong credentials fail here rather than at the first send. The check runs on the worker the mailbox is going to be placed on, so the provider sees the sign-in come from the address the mailbox will keep using. When a worker does not answer, the check moves to another one, and a row whose check no worker answered is tried again a few times before it is reported as **No worker available**.
|
||||
|
||||
Rows past the workspace's [mailbox allowance](/guides/mailboxes/#mailbox-allowance) fail as **Mailbox limit reached**; the preview says how many will fit. Raise the allowance and retry them.
|
||||
|
||||
@@ -273,6 +274,7 @@ Failures are grouped by cause, each with the fix. The same keys appear as `cause
|
||||
| `creator_removed` | Import owner left the workspace |
|
||||
| `vendor_unauthorized` | Vendor refused the API key |
|
||||
| `vendor_unreachable` | Vendor did not return the credentials |
|
||||
| `vendor_authorizing` | Vendor is authorizing Warmbly (not a failure: the row connects on its own) |
|
||||
| `mailbox_grant_not_configured` | Admin connections are not set up here |
|
||||
| `google_delegation_unauthorized` | Google refused the domain-wide delegation |
|
||||
| `microsoft_consent_missing` | Microsoft 365 admin consent is missing |
|
||||
|
||||
@@ -20,6 +20,8 @@ Open **Accounts** and choose **Add account**. Google and Microsoft each have one
|
||||
|
||||
On a self-hosted instance the two whole-domain methods are shown switched off until the operator sets them up; see [whole-domain mailbox connects](/development/whole-domain-connect/).
|
||||
|
||||
Once you press **Connect**, the connect finishes on the server even if you refresh or close the page: the credentials are checked, and the mailbox is either saved, placed on a worker and added to your list, or not saved at all. It appears in the list on its own when it is done. Pressing **Connect** again meanwhile cannot add the address twice.
|
||||
|
||||
<Callout title="Per-mailbox Google sign-in is being retired">
|
||||
A mailbox you already connected with **Sign in with Google** is not affected yet: it keeps sending and syncing, and if Google invalidates its token the **Re-authorize** button in its drawer still works. There is no cutoff date. The dashboard marks these mailboxes and offers to move each one, without losing anything, onto its domain's grant (Google Workspace) or an app password (personal Gmail). See [moving off per-mailbox Google sign-in](#moving-off-per-mailbox-google-sign-in).
|
||||
</Callout>
|
||||
|
||||
@@ -653,6 +653,14 @@ func (s *Service) Users(ctx context.Context, orgID, id uuid.UUID) ([]models.Dire
|
||||
out = append(out, models.DirectoryUser{ID: u.ID, Email: u.Mail, Name: u.DisplayName, Enabled: u.AccountEnabled})
|
||||
}
|
||||
}
|
||||
// The directory can hold domains the grant does not cover; they are not this workspace's.
|
||||
covered := out[:0]
|
||||
for _, u := range out {
|
||||
if covers(g, u.Email) {
|
||||
covered = append(covered, u)
|
||||
}
|
||||
}
|
||||
out = covered
|
||||
emails := make([]string, len(out))
|
||||
for i := range out {
|
||||
emails[i] = out[i].Email
|
||||
|
||||
@@ -202,6 +202,9 @@ func newProviders(t *testing.T) *fakeProviders {
|
||||
writeJSON(w, map[string]any{"access_token": "x", "token_type": "Bearer", "expires_in": 3600,
|
||||
"id_token": fakeIDToken(map[string]any{"tid": f.consentTenant, "wids": f.roles()})})
|
||||
})
|
||||
mux.HandleFunc("/ms/contoso.com/v2.0/.well-known/openid-configuration", func(w http.ResponseWriter, r *http.Request) {
|
||||
writeJSON(w, map[string]any{"issuer": "https://login.microsoftonline.com/11111111-1111-1111-1111-111111111111/v2.0"})
|
||||
})
|
||||
mux.HandleFunc("/ms/", func(w http.ResponseWriter, r *http.Request) {
|
||||
if !strings.Contains(r.URL.Path, "11111111-1111-1111-1111-111111111111") {
|
||||
w.WriteHeader(http.StatusBadRequest)
|
||||
|
||||
@@ -0,0 +1,134 @@
|
||||
package delegation
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"net/url"
|
||||
"strings"
|
||||
|
||||
"github.com/google/uuid"
|
||||
|
||||
"github.com/warmbly/warmbly/internal/config"
|
||||
"github.com/warmbly/warmbly/internal/errx"
|
||||
"github.com/warmbly/warmbly/internal/models"
|
||||
)
|
||||
|
||||
// microsoftAppRoles are the Graph application permissions a Microsoft grant runs on.
|
||||
var microsoftAppRoles = []string{"Mail.ReadWrite", "Mail.Send", "User.Read.All"}
|
||||
|
||||
// AppIdentity is what a vendor needs to authorize this instance's app; ok is false with no app for provider.
|
||||
func (s *Service) AppIdentity(provider string) (clientID string, scopes, appRoles []string, ok bool) {
|
||||
switch provider {
|
||||
case models.GrantProviderGoogle:
|
||||
if !s.googleEnabled() {
|
||||
return "", nil, nil, false
|
||||
}
|
||||
return s.googleClientID, append([]string(nil), config.GoogleDelegationScopes...), nil, true
|
||||
case models.GrantProviderMicrosoft:
|
||||
if !s.microsoftEnabled() {
|
||||
return "", nil, nil, false
|
||||
}
|
||||
return s.msClientID, []string{"User.Read"}, append([]string(nil), microsoftAppRoles...), true
|
||||
}
|
||||
return "", nil, nil, false
|
||||
}
|
||||
|
||||
// GrantFromVendor records a grant a vendor authorized on domain, covering only
|
||||
// domains in owned: the vendor account listing them is the workspace's proof.
|
||||
func (s *Service) GrantFromVendor(ctx context.Context, orgID, userID uuid.UUID, provider, domain, admin string, owned []string) (*models.DomainGrant, *errx.Error) {
|
||||
domain = strings.ToLower(strings.TrimSpace(domain))
|
||||
var (
|
||||
g *models.DomainGrant
|
||||
xerr *errx.Error
|
||||
)
|
||||
switch provider {
|
||||
case models.GrantProviderGoogle:
|
||||
if !s.googleEnabled() {
|
||||
return nil, notConfigured("Google Workspace")
|
||||
}
|
||||
admin = strings.ToLower(strings.TrimSpace(admin))
|
||||
if !strings.HasSuffix(admin, "@"+domain) {
|
||||
return nil, errx.NewWithIdentifier(errx.BadRequest, ErrIDProof, "The vendor has no administrator mailbox on "+domain+".")
|
||||
}
|
||||
s.forget("g\x00" + admin)
|
||||
g, xerr = s.verifyGoogle(ctx, domain, admin)
|
||||
case models.GrantProviderMicrosoft:
|
||||
if !s.microsoftEnabled() {
|
||||
return nil, notConfigured("Microsoft 365")
|
||||
}
|
||||
tenant, err := s.microsoftTenant(ctx, domain)
|
||||
if err != nil {
|
||||
return nil, unavailableErr()
|
||||
}
|
||||
s.forget("m\x00" + tenant)
|
||||
g, xerr = s.verifyMicrosoft(ctx, tenant)
|
||||
default:
|
||||
return nil, errx.ErrNotFound
|
||||
}
|
||||
if xerr != nil {
|
||||
return nil, xerr
|
||||
}
|
||||
|
||||
ownedSet := map[string]bool{}
|
||||
for _, d := range owned {
|
||||
ownedSet[strings.ToLower(strings.TrimSpace(d))] = true
|
||||
}
|
||||
domains := map[string]bool{}
|
||||
for _, d := range g.Domains {
|
||||
if ownedSet[d] {
|
||||
domains[d] = true
|
||||
}
|
||||
}
|
||||
if !domains[domain] {
|
||||
return nil, errx.NewWithIdentifier(errx.BadRequest, ErrIDNotCovered, "The authorized organization has no mailboxes on "+domain+".")
|
||||
}
|
||||
// A second domain in the same tenant widens the grant instead of replacing it.
|
||||
if list, err := s.repo.List(ctx, orgID); err == nil {
|
||||
for _, prev := range list {
|
||||
if prev.Provider == g.Provider && prev.Tenant == g.Tenant {
|
||||
for _, d := range prev.Domains {
|
||||
domains[d] = true
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
g.Domains = sortedKeys(domains)
|
||||
g.ID, g.OrganizationID = uuid.New(), orgID
|
||||
if err := s.repo.Upsert(ctx, g, userID); err != nil {
|
||||
return nil, errx.InternalError()
|
||||
}
|
||||
return s.get(ctx, orgID, g.ID)
|
||||
}
|
||||
|
||||
// microsoftTenant reads the tenant id a domain signs in to from its OpenID configuration.
|
||||
func (s *Service) microsoftTenant(ctx context.Context, domain string) (string, error) {
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodGet, s.msLoginBase+"/"+url.PathEscape(domain)+"/v2.0/.well-known/openid-configuration", nil)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
resp, err := s.http.Do(req)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
return "", fmt.Errorf("openid configuration: status %d", resp.StatusCode)
|
||||
}
|
||||
var doc struct {
|
||||
Issuer string `json:"issuer"`
|
||||
}
|
||||
if err := json.NewDecoder(io.LimitReader(resp.Body, 1<<20)).Decode(&doc); err != nil {
|
||||
return "", err
|
||||
}
|
||||
// https://login.microsoftonline.com/<tenant>/v2.0
|
||||
parts := strings.Split(strings.TrimSuffix(doc.Issuer, "/"), "/")
|
||||
for _, p := range parts {
|
||||
if id, err := uuid.Parse(p); err == nil {
|
||||
return id.String(), nil
|
||||
}
|
||||
}
|
||||
return "", fmt.Errorf("openid configuration names no tenant")
|
||||
}
|
||||
@@ -0,0 +1,74 @@
|
||||
package delegation
|
||||
|
||||
import (
|
||||
"context"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/google/uuid"
|
||||
|
||||
"github.com/warmbly/warmbly/internal/models"
|
||||
)
|
||||
|
||||
func TestVendorGrantCoversOnlyTheVendorsDomains(t *testing.T) {
|
||||
s, _, store, _, _ := newTestService(t)
|
||||
org, user := uuid.New(), uuid.New()
|
||||
|
||||
g, xerr := s.GrantFromVendor(context.Background(), org, user, models.GrantProviderMicrosoft, "Contoso.com", "", []string{"contoso.com", "elsewhere.io"})
|
||||
if xerr != nil {
|
||||
t.Fatal(xerr)
|
||||
}
|
||||
if g.Tenant != "11111111-1111-1111-1111-111111111111" || strings.Join(g.Domains, ",") != "contoso.com" {
|
||||
t.Fatalf("grant = %+v, want the tenant's domains the vendor holds and nothing else", g)
|
||||
}
|
||||
|
||||
users, xerr := s.Users(context.Background(), org, g.ID)
|
||||
if xerr != nil {
|
||||
t.Fatal(xerr)
|
||||
}
|
||||
for _, u := range users {
|
||||
if !strings.HasSuffix(u.Email, "@contoso.com") {
|
||||
t.Fatalf("directory listed %s outside the grant", u.Email)
|
||||
}
|
||||
}
|
||||
|
||||
narrow := uuid.New()
|
||||
g, xerr = s.GrantFromVendor(context.Background(), narrow, user, models.GrantProviderMicrosoft, "contoso.com", "", []string{"contoso.com"})
|
||||
if xerr != nil {
|
||||
t.Fatal(xerr)
|
||||
}
|
||||
store.grants[g.ID].Domains = []string{"contoso.onmicrosoft.com"}
|
||||
if users, _ := s.Users(context.Background(), narrow, g.ID); len(users) != 0 {
|
||||
t.Fatalf("directory listed %+v on a domain the grant does not cover", users)
|
||||
}
|
||||
|
||||
if _, xerr := s.GrantFromVendor(context.Background(), uuid.New(), user, models.GrantProviderMicrosoft, "contoso.com", "", []string{"elsewhere.io"}); xerr == nil || xerr.Identifier != ErrIDNotCovered {
|
||||
t.Fatalf("a vendor account that does not hold the domain was granted it: %v", xerr)
|
||||
}
|
||||
}
|
||||
|
||||
func TestVendorGoogleGrantNeedsTheAdminMailbox(t *testing.T) {
|
||||
s, _, _, _, _ := newTestService(t)
|
||||
if _, xerr := s.GrantFromVendor(context.Background(), uuid.New(), uuid.New(), models.GrantProviderGoogle, "acme.io", "", []string{"acme.io"}); xerr == nil || xerr.Identifier != ErrIDProof {
|
||||
t.Fatalf("a Google grant with no administrator = %v", xerr)
|
||||
}
|
||||
g, xerr := s.GrantFromVendor(context.Background(), uuid.New(), uuid.New(), models.GrantProviderGoogle, "acme.io", "admin@acme.io", []string{"acme.io"})
|
||||
if xerr != nil {
|
||||
t.Fatal(xerr)
|
||||
}
|
||||
if strings.Join(g.Domains, ",") != "acme.io" || g.AdminEmail != "admin@acme.io" {
|
||||
t.Fatalf("grant = %+v", g)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAppIdentityNamesWhatTheVendorAuthorizes(t *testing.T) {
|
||||
s, _, _, _, _ := newTestService(t)
|
||||
id, scopes, roles, ok := s.AppIdentity(models.GrantProviderMicrosoft)
|
||||
if !ok || id != "app-id" || strings.Join(roles, ",") != "Mail.ReadWrite,Mail.Send,User.Read.All" || len(scopes) == 0 {
|
||||
t.Fatalf("microsoft = %q %v %v %v", id, scopes, roles, ok)
|
||||
}
|
||||
id, scopes, roles, ok = s.AppIdentity(models.GrantProviderGoogle)
|
||||
if !ok || id != "123456789" || len(scopes) != 3 || roles != nil {
|
||||
t.Fatalf("google = %q %v %v %v", id, scopes, roles, ok)
|
||||
}
|
||||
}
|
||||
@@ -24,6 +24,8 @@ func (s *emailService) OAuthAuthorizeURL(provider models.InboxProvider, state st
|
||||
}
|
||||
|
||||
func (s *emailService) OAuthConnectWithCode(ctx context.Context, userID string, orgID *uuid.UUID, provider models.InboxProvider, code string) (*models.Email, *errx.Error) {
|
||||
ctx, cancel := detach(ctx, connectBudget)
|
||||
defer cancel()
|
||||
if code = strings.TrimSpace(code); code == "" {
|
||||
return nil, errx.ErrEmailOnboardCode
|
||||
}
|
||||
|
||||
@@ -20,6 +20,8 @@ import (
|
||||
// connect path, so a row is never connected twice and every side effect of a
|
||||
// single connect (worker load, warmup pool, webhook) happens per mailbox.
|
||||
func (s *emailService) OnboardSMTPIMAPBulk(ctx context.Context, userID string, orgID *uuid.UUID, rows []models.NewSMTPIMAPAccount) *models.MailboxBulkResult {
|
||||
ctx, cancel := detach(ctx, bulkConnectBudget)
|
||||
defer cancel()
|
||||
res := &models.MailboxBulkResult{Data: make([]models.MailboxBulkRow, len(rows))}
|
||||
res.Summary.Total = len(rows)
|
||||
if len(rows) == 0 {
|
||||
|
||||
@@ -0,0 +1,23 @@
|
||||
package email
|
||||
|
||||
import (
|
||||
"context"
|
||||
"time"
|
||||
)
|
||||
|
||||
// connectBudget bounds a single connect once it runs on its own: two silent workers, the save and the load.
|
||||
const connectBudget = 2 * time.Minute
|
||||
|
||||
// bulkConnectBudget bounds a bulk batch, which checks its rows a few at a time.
|
||||
const bulkConnectBudget = 5 * time.Minute
|
||||
|
||||
// detach keeps a connect running when the caller goes away (a refresh, a closed
|
||||
// tab), so a checked mailbox is always saved, placed, announced and loaded, or
|
||||
// not saved at all. An earlier deadline the caller set still applies.
|
||||
func detach(ctx context.Context, budget time.Duration) (context.Context, context.CancelFunc) {
|
||||
deadline := time.Now().Add(budget)
|
||||
if d, ok := ctx.Deadline(); ok && d.Before(deadline) {
|
||||
deadline = d
|
||||
}
|
||||
return context.WithDeadline(context.WithoutCancel(ctx), deadline)
|
||||
}
|
||||
@@ -0,0 +1,28 @@
|
||||
package email
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
func TestDetachOutlivesTheCallerButNotItsDeadline(t *testing.T) {
|
||||
caller, leave := context.WithCancel(context.Background())
|
||||
ctx, cancel := detach(caller, time.Minute)
|
||||
defer cancel()
|
||||
leave()
|
||||
if ctx.Err() != nil {
|
||||
t.Fatal("a refresh cancelled the connect")
|
||||
}
|
||||
if d, ok := ctx.Deadline(); !ok || time.Until(d) > time.Minute {
|
||||
t.Fatalf("deadline = %v, %v", d, ok)
|
||||
}
|
||||
|
||||
short, stop := context.WithTimeout(context.Background(), time.Second)
|
||||
defer stop()
|
||||
ctx2, cancel2 := detach(short, time.Hour)
|
||||
defer cancel2()
|
||||
if d, _ := ctx2.Deadline(); time.Until(d) > time.Second {
|
||||
t.Fatalf("the caller's earlier deadline was dropped: %v", d)
|
||||
}
|
||||
}
|
||||
@@ -68,6 +68,8 @@ func oauthMailHost(provider models.InboxProvider, email string) string {
|
||||
// and puts it to work like any other connect. The caller has already proved a
|
||||
// token can be minted for it.
|
||||
func (s *emailService) ConnectDelegated(ctx context.Context, userID string, orgID *uuid.UUID, data models.NewDelegatedAccount) (*models.Email, *errx.Error) {
|
||||
ctx, cancel := detach(ctx, connectBudget)
|
||||
defer cancel()
|
||||
if orgID == nil {
|
||||
return nil, errx.ErrNoOrganization
|
||||
}
|
||||
@@ -171,6 +173,8 @@ func (s *emailService) stoppedBySignin(ctx context.Context, accountID uuid.UUID)
|
||||
// SwitchToAppPassword moves a mailbox off per-mailbox Google sign-in onto
|
||||
// Gmail's IMAP and SMTP with an app password, keeping the mailbox and its history.
|
||||
func (s *emailService) SwitchToAppPassword(ctx context.Context, orgID *uuid.UUID, accountID uuid.UUID, appPassword string) (*models.Email, *errx.Error) {
|
||||
ctx, cancel := detach(ctx, connectBudget)
|
||||
defer cancel()
|
||||
if orgID == nil {
|
||||
return nil, errx.ErrNoOrganization
|
||||
}
|
||||
@@ -199,11 +203,7 @@ func (s *emailService) SwitchToAppPassword(ctx context.Context, orgID *uuid.UUID
|
||||
if s.workerAssignment == nil {
|
||||
return nil, errx.ErrEmailOnboardNoWorker
|
||||
}
|
||||
w, werr := s.workerAssignment.SelectValidationWorker(ctx)
|
||||
if werr != nil || w == nil {
|
||||
return nil, errx.ErrEmailOnboardNoWorker
|
||||
}
|
||||
if xerr := s.ValidateCredentials(ctx, *orgID, w.ID.String(), creds); xerr != nil {
|
||||
if xerr := s.checkCredentials(ctx, *orgID, acc.WorkerID, creds); xerr != nil {
|
||||
return nil, xerr
|
||||
}
|
||||
ok, xerr := s.emailRepository.ConvertGoogleToAppPassword(ctx, *orgID, acc.ID, creds, oauthMailHost(models.InboxProviderGoogle, acc.Email))
|
||||
|
||||
@@ -3,6 +3,7 @@ package email
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"math/rand"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
@@ -29,22 +30,26 @@ func (s *emailService) WireEmailHistoryID(repo repository.EmailHistoryIDReposito
|
||||
s.historyID = repo
|
||||
}
|
||||
|
||||
// reconcileRepublishInterval bounds how often the reconciler re-publishes a
|
||||
// given account. The immediate onboarding load and any reassignment still fire
|
||||
// right away (they call LoadAccountOntoWorker directly); this only throttles the
|
||||
// steady-state safety-net loop so the fleet isn't re-shipping every account's
|
||||
// decrypted credentials over Kafka every tick. A restarted worker is re-seeded
|
||||
// within this window rather than within one tick.
|
||||
const reconcileRepublishInterval = 5 * time.Minute
|
||||
// reconcileRepublishInterval is how often the safety net re-ships a mailbox
|
||||
// that nothing else changed. A new placement, a move and a dead worker are
|
||||
// acted on within one tick; onboarding and a worker's boot reload ship at once.
|
||||
// Each mailbox's turn is spread across the interval, because re-shipping the
|
||||
// whole fleet in one tick queued hundreds of commands on every worker at once.
|
||||
const reconcileRepublishInterval = 30 * time.Minute
|
||||
|
||||
// reconcileEntry is what the reconciler remembers about one mailbox.
|
||||
type reconcileEntry struct {
|
||||
worker uuid.UUID
|
||||
next time.Time
|
||||
}
|
||||
|
||||
// StartWorkerReconciler periodically ensures every active mailbox is assigned to
|
||||
// a worker and loaded onto it. Workers hold accounts in memory only, so this is
|
||||
// what makes onboarding, worker restarts, and reassignment converge. Each
|
||||
// account is republished at most once per reconcileRepublishInterval;
|
||||
// what makes onboarding, worker restarts, and reassignment converge.
|
||||
// PublishAddEmail is idempotent worker-side, so a republish is always safe.
|
||||
func (s *emailService) StartWorkerReconciler(ctx context.Context, interval time.Duration) {
|
||||
lastPublished := map[uuid.UUID]time.Time{}
|
||||
s.reconcileWorkerAccounts(ctx, lastPublished)
|
||||
seen := map[uuid.UUID]reconcileEntry{}
|
||||
s.reconcileWorkerAccounts(ctx, seen)
|
||||
|
||||
ticker := time.NewTicker(interval)
|
||||
defer ticker.Stop()
|
||||
@@ -53,39 +58,79 @@ func (s *emailService) StartWorkerReconciler(ctx context.Context, interval time.
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case <-ticker.C:
|
||||
s.reconcileWorkerAccounts(ctx, lastPublished)
|
||||
s.reconcileWorkerAccounts(ctx, seen)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (s *emailService) reconcileWorkerAccounts(ctx context.Context, lastPublished map[uuid.UUID]time.Time) {
|
||||
ids, err := s.emailRepository.ListActiveWorkerAccounts(ctx)
|
||||
// spreadTurn is a random point in the second half of the interval, so turns never bunch up again.
|
||||
func spreadTurn(now time.Time) time.Time {
|
||||
half := int64(reconcileRepublishInterval / 2)
|
||||
return now.Add(time.Duration(half + rand.Int63n(half)))
|
||||
}
|
||||
|
||||
func (s *emailService) reconcileWorkerAccounts(ctx context.Context, seen map[uuid.UUID]reconcileEntry) {
|
||||
rows, err := s.emailRepository.ListActiveWorkerAccounts(ctx)
|
||||
if err != nil {
|
||||
log.Warn().Err(err).Msg("worker reconciler: list active accounts failed")
|
||||
return
|
||||
}
|
||||
|
||||
active := make(map[uuid.UUID]struct{}, len(ids))
|
||||
live := map[uuid.UUID]bool{}
|
||||
isLive := func(id uuid.UUID) bool {
|
||||
v, ok := live[id]
|
||||
if !ok {
|
||||
v = true
|
||||
if s.workerAssignment != nil {
|
||||
if ok, err := s.workerAssignment.IsWorkerLive(ctx, id); err == nil {
|
||||
v = ok
|
||||
}
|
||||
}
|
||||
live[id] = v
|
||||
}
|
||||
return v
|
||||
}
|
||||
now := time.Now()
|
||||
for _, id := range ids {
|
||||
active[id] = struct{}{}
|
||||
if last, ok := lastPublished[id]; ok && now.Sub(last) < reconcileRepublishInterval {
|
||||
for _, r := range reconcileDue(rows, seen, now, isLive) {
|
||||
if err := s.LoadAccountOntoWorker(ctx, r.ID); err != nil {
|
||||
log.Warn().Err(err).Str("email_id", r.ID.String()).Msg("worker reconciler: load account failed")
|
||||
continue
|
||||
}
|
||||
if err := s.LoadAccountOntoWorker(ctx, id); err != nil {
|
||||
log.Warn().Err(err).Str("email_id", id.String()).Msg("worker reconciler: load account failed")
|
||||
continue
|
||||
var worker uuid.UUID
|
||||
if r.WorkerID != nil {
|
||||
worker = *r.WorkerID
|
||||
}
|
||||
lastPublished[id] = now
|
||||
seen[r.ID] = reconcileEntry{worker: worker, next: spreadTurn(now)}
|
||||
}
|
||||
}
|
||||
|
||||
// Drop throttle entries for accounts no longer active so the map can't grow
|
||||
// without bound as mailboxes are disconnected.
|
||||
for id := range lastPublished {
|
||||
// reconcileDue picks the mailboxes to ship this tick: at once when unplaced,
|
||||
// on a dead worker or moved, otherwise when their spread-out turn comes. It
|
||||
// also schedules first sightings and forgets mailboxes that are gone.
|
||||
func reconcileDue(rows []repository.MailboxAssignment, seen map[uuid.UUID]reconcileEntry, now time.Time, isLive func(uuid.UUID) bool) []repository.MailboxAssignment {
|
||||
var due []repository.MailboxAssignment
|
||||
active := make(map[uuid.UUID]struct{}, len(rows))
|
||||
for _, r := range rows {
|
||||
active[r.ID] = struct{}{}
|
||||
entry, known := seen[r.ID]
|
||||
urgent := r.WorkerID == nil || !isLive(*r.WorkerID) || (known && entry.worker != *r.WorkerID)
|
||||
switch {
|
||||
case urgent:
|
||||
case !known:
|
||||
// First sight since boot: onboarding or the worker's boot reload
|
||||
// already shipped it, so its safety-net turn is spread out.
|
||||
seen[r.ID] = reconcileEntry{worker: *r.WorkerID, next: now.Add(time.Duration(rand.Int63n(int64(reconcileRepublishInterval))))}
|
||||
continue
|
||||
case now.Before(entry.next):
|
||||
continue
|
||||
}
|
||||
due = append(due, r)
|
||||
}
|
||||
for id := range seen {
|
||||
if _, ok := active[id]; !ok {
|
||||
delete(lastPublished, id)
|
||||
delete(seen, id)
|
||||
}
|
||||
}
|
||||
return due
|
||||
}
|
||||
|
||||
// ReloadWorkerAccounts publishes every active mailbox assigned to workerID
|
||||
|
||||
@@ -107,6 +107,8 @@ func (s *emailService) guardInboxLimit(ctx context.Context, orgID *uuid.UUID) (*
|
||||
// inbox owner, and persists a new email account — or, when the state carries an
|
||||
// account id (OAuthReauth), renews that mailbox's tokens in place instead.
|
||||
func (s *emailService) OAuthFinish(ctx context.Context, userID, code, state string) (*models.Email, bool, *errx.Error) {
|
||||
ctx, cancel := detach(ctx, connectBudget)
|
||||
defer cancel()
|
||||
if code = strings.TrimSpace(code); code == "" {
|
||||
return nil, false, errx.ErrEmailOnboardCode
|
||||
}
|
||||
@@ -231,6 +233,8 @@ func (s *emailService) OAuthFinish(ctx context.Context, userID, code, state stri
|
||||
// OnboardSMTPIMAP validates the supplied SMTP/IMAP credentials against a live worker, then
|
||||
// persists the email account on success. Returns ErrEmailCredentials if the worker reports failure.
|
||||
func (s *emailService) OnboardSMTPIMAP(ctx context.Context, userID string, orgID *uuid.UUID, data *models.NewSMTPIMAPAccount) (*models.Email, *errx.Error) {
|
||||
ctx, cancel := detach(ctx, connectBudget)
|
||||
defer cancel()
|
||||
if xerr := validateSMTPIMAPInput(data); xerr != nil {
|
||||
return nil, xerr
|
||||
}
|
||||
@@ -250,15 +254,8 @@ func (s *emailService) OnboardSMTPIMAP(ctx context.Context, userID string, orgID
|
||||
return nil, errx.ErrEmailOnboardNoWorker
|
||||
}
|
||||
|
||||
// Any live worker can run the one-shot validation handshake: nothing is
|
||||
// placed yet, the worker just dials the credentials once and reports back.
|
||||
w, werr := s.workerAssignment.SelectValidationWorker(ctx)
|
||||
if werr != nil || w == nil {
|
||||
return nil, errx.ErrEmailOnboardNoWorker
|
||||
}
|
||||
|
||||
creds := &models.SmtpImap{SMTP: data.SMTP, IMAP: data.IMAP}
|
||||
if xerr := s.ValidateCredentials(ctx, *orgID, w.ID.String(), creds); xerr != nil {
|
||||
if xerr := s.checkCredentials(ctx, *orgID, nil, creds); xerr != nil {
|
||||
return nil, xerr
|
||||
}
|
||||
|
||||
|
||||
@@ -134,6 +134,8 @@ func (s *emailService) finishReauth(ctx context.Context, sess *models.EmailOnboa
|
||||
// validate the replacement credentials against a live worker, store them, and
|
||||
// put the mailbox back to work.
|
||||
func (s *emailService) UpdateSMTPIMAPCredentials(ctx context.Context, orgID *uuid.UUID, accountID uuid.UUID, creds *models.SmtpImap) (*models.Email, *errx.Error) {
|
||||
ctx, cancel := detach(ctx, connectBudget)
|
||||
defer cancel()
|
||||
if orgID == nil {
|
||||
return nil, errx.ErrNoOrganization
|
||||
}
|
||||
@@ -158,13 +160,7 @@ func (s *emailService) UpdateSMTPIMAPCredentials(ctx context.Context, orgID *uui
|
||||
if s.workerAssignment == nil {
|
||||
return nil, errx.ErrEmailOnboardNoWorker
|
||||
}
|
||||
// Any live worker can run the one-shot validation handshake, same as at
|
||||
// connect time.
|
||||
w, werr := s.workerAssignment.SelectValidationWorker(ctx)
|
||||
if werr != nil || w == nil {
|
||||
return nil, errx.ErrEmailOnboardNoWorker
|
||||
}
|
||||
if xerr := s.ValidateCredentials(ctx, *orgID, w.ID.String(), creds); xerr != nil {
|
||||
if xerr := s.checkCredentials(ctx, *orgID, account.WorkerID, creds); xerr != nil {
|
||||
return nil, xerr
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,67 @@
|
||||
package email
|
||||
|
||||
import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
|
||||
"github.com/warmbly/warmbly/internal/repository"
|
||||
)
|
||||
|
||||
func dueIDs(rows []repository.MailboxAssignment) map[uuid.UUID]bool {
|
||||
out := map[uuid.UUID]bool{}
|
||||
for _, r := range rows {
|
||||
out[r.ID] = true
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
func TestReconcileShipsChangesAtOnceAndSpreadsTheRest(t *testing.T) {
|
||||
w1, w2, dead := uuid.New(), uuid.New(), uuid.New()
|
||||
live := func(id uuid.UUID) bool { return id != dead }
|
||||
var steady []repository.MailboxAssignment
|
||||
for i := 0; i < 200; i++ {
|
||||
steady = append(steady, repository.MailboxAssignment{ID: uuid.New(), WorkerID: &w1})
|
||||
}
|
||||
unplaced := repository.MailboxAssignment{ID: uuid.New()}
|
||||
onDead := repository.MailboxAssignment{ID: uuid.New(), WorkerID: &dead}
|
||||
rows := append(append([]repository.MailboxAssignment(nil), steady...), unplaced, onDead)
|
||||
|
||||
seen := map[uuid.UUID]reconcileEntry{}
|
||||
now := time.Now()
|
||||
due := dueIDs(reconcileDue(rows, seen, now, live))
|
||||
if len(due) != 2 || !due[unplaced.ID] || !due[onDead.ID] {
|
||||
t.Fatalf("first pass shipped %d mailboxes, want only the unplaced one and the one on a dead worker", len(due))
|
||||
}
|
||||
|
||||
// A moved mailbox ships on the next tick, not at its turn.
|
||||
seen[steady[0].ID] = reconcileEntry{worker: w1, next: now.Add(time.Hour)}
|
||||
steady[0].WorkerID = &w2
|
||||
rows[0] = steady[0]
|
||||
if due := dueIDs(reconcileDue(rows, seen, now.Add(time.Minute), live)); !due[steady[0].ID] {
|
||||
t.Fatal("a moved mailbox waited for its turn")
|
||||
}
|
||||
|
||||
// Over one interval every steady mailbox gets exactly one turn, never all in one tick.
|
||||
shipped, worst := map[uuid.UUID]int{}, 0
|
||||
for tick := now; !tick.After(now.Add(reconcileRepublishInterval)); tick = tick.Add(time.Minute) {
|
||||
due := reconcileDue(steady[1:], seen, tick, live)
|
||||
if len(due) > worst {
|
||||
worst = len(due)
|
||||
}
|
||||
for _, r := range due {
|
||||
shipped[r.ID]++
|
||||
seen[r.ID] = reconcileEntry{worker: w1, next: tick.Add(reconcileRepublishInterval)}
|
||||
}
|
||||
}
|
||||
if len(shipped) != len(steady)-1 || worst > 40 {
|
||||
t.Fatalf("shipped %d of %d, worst tick %d: turns must cover everyone and stay spread", len(shipped), len(steady)-1, worst)
|
||||
}
|
||||
|
||||
// A mailbox that is no longer active is forgotten.
|
||||
reconcileDue(nil, seen, now, live)
|
||||
if len(seen) != 0 {
|
||||
t.Fatalf("%d inactive mailboxes remembered", len(seen))
|
||||
}
|
||||
}
|
||||
@@ -3,88 +3,157 @@ package email
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/redis/go-redis/v9"
|
||||
"github.com/rs/zerolog/log"
|
||||
"github.com/warmbly/warmbly/internal/errx"
|
||||
"github.com/warmbly/warmbly/internal/models"
|
||||
"github.com/warmbly/warmbly/internal/observability/errs"
|
||||
)
|
||||
|
||||
// ValidateCredentials seals a copy of the credentials with the org DEK and asks
|
||||
// a worker to try them against the live servers. The caller's credentials are
|
||||
// never mutated, except for the SMTP port when the worker signed in on another
|
||||
// one (adoptProbedPort): they go on to be stored under the credentials key, and
|
||||
// sealing them in place here would double-encrypt the stored password.
|
||||
func (s *emailService) ValidateCredentials(ctx context.Context, orgID uuid.UUID, workerID string, credentials *models.SmtpImap) *errx.Error {
|
||||
processID := uuid.New()
|
||||
// checkCredentials runs a credential check on the best worker that takes it.
|
||||
// prefer is the mailbox's own worker on a reconnect, nil on a first connect.
|
||||
func (s *emailService) checkCredentials(ctx context.Context, orgID uuid.UUID, prefer *uuid.UUID, creds *models.SmtpImap) *errx.Error {
|
||||
if s.workerAssignment == nil {
|
||||
return errx.ErrEmailOnboardNoWorker
|
||||
}
|
||||
workers, err := s.workerAssignment.ValidationWorkers(ctx, orgID, prefer)
|
||||
if err != nil || len(workers) == 0 {
|
||||
return errx.ErrEmailOnboardNoWorker
|
||||
}
|
||||
return s.ValidateCredentials(ctx, orgID, workers, creds)
|
||||
}
|
||||
|
||||
// ValidateCredentials asks the workers in order to try the credentials, over
|
||||
// each one's Redis request channel so no queued command delays the check. A
|
||||
// worker not listening is skipped, a silent one costs one wait, and the bus is
|
||||
// used only when none listens. Only the SMTP port may be rewritten (adoptProbedPort).
|
||||
func (s *emailService) ValidateCredentials(ctx context.Context, orgID uuid.UUID, workers []models.Worker, credentials *models.SmtpImap) *errx.Error {
|
||||
if credentials == nil || credentials.SMTP == nil || credentials.IMAP == nil {
|
||||
return errx.ErrEmailCredentialsRequired
|
||||
}
|
||||
if len(workers) == 0 {
|
||||
return errx.ErrEmailOnboardNoWorker
|
||||
}
|
||||
|
||||
cipher, err := s.cipherService.Cipher(ctx, orgID)
|
||||
if err != nil {
|
||||
errs.CaptureException(err)
|
||||
return errx.InternalError()
|
||||
}
|
||||
|
||||
sealedIMAP := *credentials.IMAP
|
||||
sealedIMAP.Password, err = cipher.Encrypt(ctx, credentials.IMAP.Password)
|
||||
if err != nil {
|
||||
errs.CaptureException(err)
|
||||
return errx.InternalError()
|
||||
}
|
||||
|
||||
sealedSMTP := *credentials.SMTP
|
||||
sealedSMTP.Password, err = cipher.Encrypt(ctx, credentials.SMTP.Password)
|
||||
if err != nil {
|
||||
errs.CaptureException(err)
|
||||
return errx.InternalError()
|
||||
}
|
||||
sealed := &models.SmtpImap{SMTP: &sealedSMTP, IMAP: &sealedIMAP}
|
||||
|
||||
// This side must wait longer than the worker's own budget, or the answer
|
||||
// arrives after the only listener has given up and every slow mail host
|
||||
// reads as an outage.
|
||||
subscribeContext, cancel := context.WithTimeout(ctx, validationWait)
|
||||
silent := 0
|
||||
for _, w := range workers {
|
||||
if silent >= validationAttempts || ctx.Err() != nil {
|
||||
break
|
||||
}
|
||||
xerr, outcome := s.validateOn(ctx, orgID, w.ID, sealed, credentials, false)
|
||||
switch outcome {
|
||||
case validationAnswered:
|
||||
return xerr
|
||||
case validationSilent:
|
||||
silent++
|
||||
}
|
||||
}
|
||||
if silent > 0 || ctx.Err() != nil {
|
||||
return errx.ErrEmailValidation
|
||||
}
|
||||
xerr, outcome := s.validateOn(ctx, orgID, workers[0].ID, sealed, credentials, true)
|
||||
if outcome == validationSilent {
|
||||
return errx.ErrEmailValidation
|
||||
}
|
||||
return xerr
|
||||
}
|
||||
|
||||
type validationOutcome int
|
||||
|
||||
const (
|
||||
validationAnswered validationOutcome = iota
|
||||
// validationNotListening: nobody took the request, so nothing was tried.
|
||||
validationNotListening
|
||||
// validationSilent: a worker took the request and never answered.
|
||||
validationSilent
|
||||
)
|
||||
|
||||
// validationAttempts bounds how many silent workers one check waits on.
|
||||
const validationAttempts = 2
|
||||
|
||||
// validateOn sends one check to one worker and waits for its verdict.
|
||||
func (s *emailService) validateOn(ctx context.Context, orgID, workerID uuid.UUID, sealed, credentials *models.SmtpImap, overBus bool) (*errx.Error, validationOutcome) {
|
||||
processID := uuid.New()
|
||||
|
||||
// Longer than the worker's own budget, so a verdict at its limit is still heard.
|
||||
waitCtx, cancel := context.WithTimeout(ctx, validationWait)
|
||||
defer cancel()
|
||||
|
||||
// Subscribe before the job goes out. Redis pub/sub keeps nothing for a
|
||||
// channel with no subscriber, so a worker that answers between the publish
|
||||
// and the SUBSCRIBE lands its reply nowhere and the wait below runs to its
|
||||
// deadline with the validation already done.
|
||||
r := s.r.Subscribe(subscribeContext, "email_validation:"+processID.String())
|
||||
defer r.Close()
|
||||
|
||||
if err := s.publisher.PublishEmailValidation(ctx, workerID, models.EventWorkerEmailValidation{
|
||||
OrgID: orgID,
|
||||
ProcessID: processID,
|
||||
Credentials: &models.SmtpImap{SMTP: &sealedSMTP, IMAP: &sealedIMAP},
|
||||
}); err != nil {
|
||||
// Confirmed before the request goes out: Redis keeps nothing for a channel nobody holds.
|
||||
sub := s.r.Subscribe(waitCtx, models.EmailValidationReplyChannel(processID))
|
||||
defer sub.Close()
|
||||
if _, err := sub.Receive(waitCtx); err != nil {
|
||||
if waitCtx.Err() != nil {
|
||||
return nil, validationSilent
|
||||
}
|
||||
errs.CaptureException(err)
|
||||
return errx.InternalError()
|
||||
return errx.InternalError(), validationAnswered
|
||||
}
|
||||
|
||||
for {
|
||||
msg, err := r.ReceiveMessage(subscribeContext)
|
||||
req := models.EventWorkerEmailValidation{OrgID: orgID, ProcessID: processID, Credentials: sealed}
|
||||
if overBus {
|
||||
if err := s.publisher.PublishEmailValidation(ctx, workerID.String(), req); err != nil {
|
||||
errs.CaptureException(err)
|
||||
return errx.InternalError(), validationAnswered
|
||||
}
|
||||
} else {
|
||||
body, err := json.Marshal(req)
|
||||
if err != nil {
|
||||
// Ask the context, not the error, and ask it for the deadline
|
||||
// specifically. A deadline reached while waiting is the mail host
|
||||
// being slow, not this service being broken, but go-redis pushes
|
||||
// the deadline down onto the socket and it comes back as a net
|
||||
// timeout rather than context.DeadlineExceeded, so the error's own
|
||||
// shape cannot say whose deadline it was. The context can, and it
|
||||
// also distinguishes the two ways it ends: only the timeout below
|
||||
// is the mail host. A socket timeout while the context is still
|
||||
// live is Redis failing, and a caller who went away is neither.
|
||||
if errors.Is(subscribeContext.Err(), context.DeadlineExceeded) {
|
||||
return errx.ErrEmailValidation
|
||||
return errx.InternalError(), validationAnswered
|
||||
}
|
||||
receivers, err := s.r.Publish(ctx, models.EmailValidationRequestChannel(workerID), body).Result()
|
||||
if err != nil {
|
||||
errs.CaptureException(err)
|
||||
return errx.InternalError(), validationAnswered
|
||||
}
|
||||
if receivers == 0 {
|
||||
return nil, validationNotListening
|
||||
}
|
||||
}
|
||||
|
||||
xerr, answered := awaitVerdict(waitCtx, sub, workerID, processID, credentials)
|
||||
if !answered {
|
||||
log.Warn().Str("worker_id", workerID.String()).Bool("over_bus", overBus).Msg("mailbox validation: worker took the check and did not answer")
|
||||
return nil, validationSilent
|
||||
}
|
||||
return xerr, validationAnswered
|
||||
}
|
||||
|
||||
// awaitVerdict reads the worker's answer; answered is false when none came in time.
|
||||
func awaitVerdict(ctx context.Context, sub *redis.PubSub, workerID, processID uuid.UUID, credentials *models.SmtpImap) (*errx.Error, bool) {
|
||||
for {
|
||||
msg, err := sub.ReceiveMessage(ctx)
|
||||
if err != nil {
|
||||
// Ask the context, not the error: go-redis pushes the deadline onto
|
||||
// the socket as a net timeout. One while the context is live is Redis failing.
|
||||
if ctx.Err() != nil {
|
||||
return nil, false
|
||||
}
|
||||
errs.CaptureException(err)
|
||||
return errx.InternalError()
|
||||
return errx.InternalError(), true
|
||||
}
|
||||
|
||||
// A worker answers with a JSON verdict followed by the legacy digit.
|
||||
@@ -92,9 +161,9 @@ func (s *emailService) ValidateCredentials(ctx context.Context, orgID uuid.UUID,
|
||||
// the digit, and a payload that is neither is skipped.
|
||||
switch msg.Payload {
|
||||
case "1":
|
||||
return nil
|
||||
return nil, true
|
||||
case "0":
|
||||
return errx.ErrEmailCredentials
|
||||
return errx.ErrEmailCredentials, true
|
||||
}
|
||||
var verdict models.EmailValidationVerdict
|
||||
if err := json.Unmarshal([]byte(msg.Payload), &verdict); err != nil {
|
||||
@@ -104,15 +173,15 @@ func (s *emailService) ValidateCredentials(ctx context.Context, orgID uuid.UUID,
|
||||
// The worker answered but never reached the mail server; that is
|
||||
// this side's failure, not the customer's host or port.
|
||||
err := fmt.Errorf("mailbox validation on worker %s did not run: %s", workerID, verdict.Error)
|
||||
log.Error().Str("worker_id", workerID).Str("process_id", processID.String()).Msg(err.Error())
|
||||
log.Error().Str("worker_id", workerID.String()).Str("process_id", processID.String()).Msg(err.Error())
|
||||
errs.CaptureException(err)
|
||||
return errx.InternalError()
|
||||
return errx.InternalError(), true
|
||||
}
|
||||
if verdict.OK {
|
||||
adoptProbedPort(credentials.SMTP, verdict.SMTP)
|
||||
return nil
|
||||
return nil, true
|
||||
}
|
||||
return validationError(verdict, credentials)
|
||||
return validationError(verdict, credentials), true
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -3,6 +3,7 @@ package mailboximport
|
||||
import (
|
||||
"github.com/warmbly/warmbly/internal/models"
|
||||
"github.com/warmbly/warmbly/internal/pkg/mailcause"
|
||||
"github.com/warmbly/warmbly/internal/repository"
|
||||
)
|
||||
|
||||
// Causes an import finds before it dials anything, or around the dial.
|
||||
@@ -23,6 +24,8 @@ const (
|
||||
|
||||
causeVendorUnauthorized = "vendor_unauthorized"
|
||||
causeVendorUnreachable = "vendor_unreachable"
|
||||
// causeVendorAuthorizing parks a row while its vendor authorizes Warmbly on the domain.
|
||||
causeVendorAuthorizing = repository.ParkedVendorCause
|
||||
)
|
||||
|
||||
// ErrIDVendorUnauthorized is the identifier a VendorSource answers with when the vendor refused the key.
|
||||
@@ -57,6 +60,8 @@ var importCauses = map[string]mailcause.Cause{
|
||||
Fix: "The inbox vendor no longer accepts the API key saved for this account. Update the key from Add account > Inbox vendor, then retry these rows.", Retryable: true},
|
||||
causeVendorUnreachable: {Key: causeVendorUnreachable, Title: "Vendor did not return the credentials",
|
||||
Fix: "The inbox vendor did not answer with this mailbox's credentials, or no longer has it. Retry in a few minutes, or check the mailbox in the vendor's dashboard.", Retryable: true},
|
||||
causeVendorAuthorizing: {Key: causeVendorAuthorizing, Title: "Vendor is authorizing Warmbly",
|
||||
Fix: "The inbox vendor is approving Warmbly for these domains through the admin mailbox it holds, so nobody has to sign in. This usually takes a few minutes; the rows connect on their own when it finishes."},
|
||||
"mailbox_grant_not_configured": {Key: "mailbox_grant_not_configured", Title: "Admin connections are not set up here",
|
||||
Fix: "This instance has no Google service account or Microsoft app for admin connections. An operator configures them; see the self-hosting docs."},
|
||||
"google_delegation_unauthorized": {Key: "google_delegation_unauthorized", Title: "Google refused the domain-wide delegation",
|
||||
|
||||
@@ -0,0 +1,48 @@
|
||||
package mailboximport
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/warmbly/warmbly/internal/errx"
|
||||
)
|
||||
|
||||
func TestRetryUnansweredRetriesOnlyASilentFleet(t *testing.T) {
|
||||
unansweredBackoff, unansweredReserve = time.Millisecond, 0
|
||||
t.Cleanup(func() { unansweredBackoff, unansweredReserve = 5*time.Second, 48*time.Second })
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
||||
defer cancel()
|
||||
|
||||
calls := 0
|
||||
xerr := retryUnanswered(ctx, func() *errx.Error {
|
||||
calls++
|
||||
if calls < 3 {
|
||||
return errx.ErrEmailValidation
|
||||
}
|
||||
return nil
|
||||
})
|
||||
if xerr != nil || calls != 3 {
|
||||
t.Fatalf("silent workers: err = %v after %d calls, want success on the third", xerr, calls)
|
||||
}
|
||||
|
||||
calls = 0
|
||||
refused := errx.NewWithIdentifier(errx.BadRequest, errx.ErrEmailValidation.Identifier, "IMAP did not answer")
|
||||
refused.Cause = "imap_timeout"
|
||||
if xerr := retryUnanswered(ctx, func() *errx.Error { calls++; return refused }); xerr != refused || calls != 1 {
|
||||
t.Fatalf("a verdict about the mail server was retried: %d calls", calls)
|
||||
}
|
||||
|
||||
calls = 0
|
||||
if xerr := retryUnanswered(ctx, func() *errx.Error { calls++; return errx.ErrEmailCredentials }); xerr != errx.ErrEmailCredentials || calls != 1 {
|
||||
t.Fatalf("refused credentials were retried: %d calls", calls)
|
||||
}
|
||||
|
||||
calls = 0
|
||||
short, cancelShort := context.WithTimeout(context.Background(), time.Millisecond)
|
||||
defer cancelShort()
|
||||
unansweredReserve = time.Hour
|
||||
if xerr := retryUnanswered(short, func() *errx.Error { calls++; return errx.ErrEmailOnboardNoWorker }); xerr != errx.ErrEmailOnboardNoWorker || calls != 1 {
|
||||
t.Fatalf("retried past the lease: %d calls", calls)
|
||||
}
|
||||
}
|
||||
@@ -39,6 +39,7 @@ func (s *Service) Kick() {
|
||||
func (s *Service) Start(ctx context.Context) {
|
||||
go jobrun.Loop(ctx, "mailbox_import_upkeep", time.Minute, true, func(ctx context.Context) error {
|
||||
s.reconcileSignins(ctx)
|
||||
s.resumeVendorAuthorizations(ctx)
|
||||
s.classifyHosts(ctx)
|
||||
if err := s.repo.SettleOrphans(ctx); err != nil {
|
||||
return err
|
||||
@@ -162,6 +163,17 @@ func (s *Service) process(ctx context.Context, w repository.ImportWorkRow) {
|
||||
return
|
||||
}
|
||||
|
||||
// A restart between creating this mailbox and recording the row leaves it for a
|
||||
// later claim: finish the row's work rather than treat the mailbox as pre-existing.
|
||||
if existing != nil && w.Attempts > 1 {
|
||||
acc, xerr := s.mailboxes.Get(rowCtx, orgID.String(), existing.ID.String())
|
||||
ours := w.AccountID != nil && existing.ID == *w.AccountID
|
||||
if xerr == nil && acc != nil && (ours || acc.CreatedAt.After(w.ImportCreatedAt)) {
|
||||
s.connected(ctx, w, p, acc, settings, userID)
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
if existing != nil && w.OnExisting == "skip" {
|
||||
s.finish(ctx, w, models.ImportRowSkipped, "already_connected", "", "Already in this workspace.", &existing.ID, false)
|
||||
return
|
||||
@@ -178,6 +190,21 @@ func (s *Service) process(ctx context.Context, w repository.ImportWorkRow) {
|
||||
// A vendor row learns its credentials only now, and is judged like a file row.
|
||||
if p.VendorConnectionID != nil && p.SMTP == nil && !p.Signin && p.GrantID == nil {
|
||||
resolved, cause, problem := s.resolveVendorRow(rowCtx, w, p)
|
||||
if cause == causeMicrosoftSignin || cause == causeGoogleSignin {
|
||||
auth := s.authorizeVendorDomain(rowCtx, w, *p.VendorConnectionID, cause)
|
||||
switch {
|
||||
case auth.GrantID != nil:
|
||||
resolved.GrantID, resolved.AuthMethod, resolved.Signin = auth.GrantID, models.MailAuthDelegated, false
|
||||
resolved.SMTP, resolved.IMAP = nil, nil
|
||||
cause = ""
|
||||
case auth.Pending:
|
||||
s.finish(ctx, w, models.ImportRowNeedsSignin, cause, causeVendorAuthorizing,
|
||||
"Your inbox vendor is authorizing Warmbly on "+domainOf(w.Email)+". This mailbox connects on its own when it finishes.", nil, true)
|
||||
return
|
||||
case auth.Message != "":
|
||||
problem = auth.Message
|
||||
}
|
||||
}
|
||||
if cause != "" {
|
||||
status := models.ImportRowFailed
|
||||
if cause == causeMicrosoftSignin || cause == causeGoogleSignin {
|
||||
@@ -236,7 +263,10 @@ func (s *Service) process(ctx context.Context, w repository.ImportWorkRow) {
|
||||
return
|
||||
}
|
||||
creds := &models.SmtpImap{SMTP: p.SMTP, IMAP: p.IMAP}
|
||||
if _, xerr := s.emails.UpdateSMTPIMAPCredentials(rowCtx, &orgID, existing.ID, creds); xerr != nil {
|
||||
if xerr := retryUnanswered(rowCtx, func() *errx.Error {
|
||||
_, xerr := s.emails.UpdateSMTPIMAPCredentials(rowCtx, &orgID, existing.ID, creds)
|
||||
return xerr
|
||||
}); xerr != nil {
|
||||
s.fail(ctx, w, xerr)
|
||||
return
|
||||
}
|
||||
@@ -246,8 +276,13 @@ func (s *Service) process(ctx context.Context, w repository.ImportWorkRow) {
|
||||
return
|
||||
}
|
||||
|
||||
acc, xerr := s.emails.OnboardSMTPIMAP(rowCtx, userID, &orgID, &models.NewSMTPIMAPAccount{
|
||||
Email: w.Email, Name: p.Name, SMTP: p.SMTP, IMAP: p.IMAP, MailHost: p.MailHost, AuthMethod: p.AuthMethod,
|
||||
var acc *models.Email
|
||||
xerr = retryUnanswered(rowCtx, func() *errx.Error {
|
||||
var xerr *errx.Error
|
||||
acc, xerr = s.emails.OnboardSMTPIMAP(rowCtx, userID, &orgID, &models.NewSMTPIMAPAccount{
|
||||
Email: w.Email, Name: p.Name, SMTP: p.SMTP, IMAP: p.IMAP, MailHost: p.MailHost, AuthMethod: p.AuthMethod,
|
||||
})
|
||||
return xerr
|
||||
})
|
||||
if xerr != nil {
|
||||
if errors.Is(xerr, errx.ErrEmailOnboardAlreadyExists) {
|
||||
@@ -261,8 +296,43 @@ func (s *Service) process(ctx context.Context, w repository.ImportWorkRow) {
|
||||
s.connected(ctx, w, p, acc, settings, userID)
|
||||
}
|
||||
|
||||
// unansweredBackoff is the first wait before a connect no worker answered is tried again; it triples.
|
||||
var unansweredBackoff = 5 * time.Second
|
||||
|
||||
// unansweredReserve is what one more connect needs of the row's lease: two silent workers, then saving the mailbox.
|
||||
var unansweredReserve = 48 * time.Second
|
||||
|
||||
// retryUnanswered repeats a connect whose check no worker answered, while the row's lease leaves room.
|
||||
func retryUnanswered(ctx context.Context, connect func() *errx.Error) *errx.Error {
|
||||
for wait := unansweredBackoff; ; wait *= 3 {
|
||||
xerr := connect()
|
||||
if xerr == nil || !unanswered(xerr) {
|
||||
return xerr
|
||||
}
|
||||
if d, ok := ctx.Deadline(); ok && time.Until(d) < wait+unansweredReserve {
|
||||
return xerr
|
||||
}
|
||||
t := time.NewTimer(wait)
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
t.Stop()
|
||||
return xerr
|
||||
case <-t.C:
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// unanswered is a connect that never reached the mail server because no worker took or answered the check.
|
||||
func unanswered(xerr *errx.Error) bool {
|
||||
return errors.Is(xerr, errx.ErrEmailOnboardNoWorker) || (xerr.Cause == "" && xerr.Identifier == errx.ErrEmailValidation.Identifier)
|
||||
}
|
||||
|
||||
// connected records a new mailbox: audit, vendor link, settings, row outcome.
|
||||
func (s *Service) connected(ctx context.Context, w repository.ImportWorkRow, p payload, acc *models.Email, settings models.MailboxImportSettings, userID string) {
|
||||
// Recorded first, so a restart before the row is finished resumes here instead of skipping the mailbox.
|
||||
if w.AccountID == nil || *w.AccountID != acc.ID {
|
||||
_ = s.repo.SetRowAccount(ctx, w.ImportID, w.Line, w.Attempts, acc.ID)
|
||||
}
|
||||
if s.auditor != nil {
|
||||
s.auditor.LogAction(ctx, w.OrgID, *w.CreatedBy, models.AuditActionConnect, models.AuditEntityEmailAccount, &acc.ID, "", "", nil,
|
||||
map[string]string{"provider": acc.Provider, "email": acc.Email, "import_id": w.ImportID.String()})
|
||||
@@ -308,6 +378,65 @@ func (s *Service) link(ctx context.Context, p payload, accountID uuid.UUID) {
|
||||
}
|
||||
}
|
||||
|
||||
// authorizeVendorDomain asks the row's vendor to authorize Warmbly on its domain; signinCause names the provider.
|
||||
func (s *Service) authorizeVendorDomain(ctx context.Context, w repository.ImportWorkRow, connectionID uuid.UUID, signinCause string) VendorAuthorization {
|
||||
az, ok := s.vendors.(VendorAuthorizer)
|
||||
if !ok || w.CreatedBy == nil {
|
||||
return VendorAuthorization{}
|
||||
}
|
||||
provider := models.GrantProviderMicrosoft
|
||||
if signinCause == causeGoogleSignin {
|
||||
provider = models.GrantProviderGoogle
|
||||
}
|
||||
return az.AuthorizeDomain(ctx, w.OrgID, *w.CreatedBy, connectionID, w.Email, provider)
|
||||
}
|
||||
|
||||
// resumeVendorAuthorizations requeues rows parked on a vendor authorization once
|
||||
// it has an answer, asking the vendor once per domain.
|
||||
func (s *Service) resumeVendorAuthorizations(ctx context.Context) {
|
||||
if s.vendors == nil {
|
||||
return
|
||||
}
|
||||
rows, err := s.repo.ParkedRows(ctx, causeVendorAuthorizing, 500)
|
||||
if err != nil || len(rows) == 0 {
|
||||
return
|
||||
}
|
||||
type domainKey struct {
|
||||
org, conn uuid.UUID
|
||||
signin, domain string
|
||||
}
|
||||
settled := map[domainKey]bool{}
|
||||
resumed := false
|
||||
var waiting []repository.ImportWorkRow
|
||||
for _, w := range rows {
|
||||
ciph, err := s.cipher.Cipher(ctx, w.OrgID)
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
p, xerr := unseal(ctx, ciph, w.Payload)
|
||||
if xerr != nil || p.VendorConnectionID == nil {
|
||||
continue
|
||||
}
|
||||
k := domainKey{w.OrgID, *p.VendorConnectionID, w.Code, domainOf(w.Email)}
|
||||
done, seen := settled[k]
|
||||
if !seen {
|
||||
done = !s.authorizeVendorDomain(ctx, w, *p.VendorConnectionID, w.Code).Pending
|
||||
settled[k] = done
|
||||
}
|
||||
if !done {
|
||||
waiting = append(waiting, w)
|
||||
continue
|
||||
}
|
||||
if err := s.repo.ResumeParked(ctx, w.ImportID, w.Line, causeVendorAuthorizing); err == nil {
|
||||
resumed = true
|
||||
}
|
||||
}
|
||||
_ = s.repo.TouchParked(ctx, causeVendorAuthorizing, waiting)
|
||||
if resumed {
|
||||
s.Kick()
|
||||
}
|
||||
}
|
||||
|
||||
// resolveVendorRow fetches a vendor row's credentials and builds it the way a
|
||||
// file row is built, so detection, grants and validation all apply.
|
||||
func (s *Service) resolveVendorRow(ctx context.Context, w repository.ImportWorkRow, p payload) (payload, string, string) {
|
||||
|
||||
@@ -85,6 +85,21 @@ type VendorSource interface {
|
||||
Link(ctx context.Context, accountID, connectionID uuid.UUID, mailboxID string)
|
||||
}
|
||||
|
||||
// VendorAuthorization is where getting a vendor domain onto a grant stands.
|
||||
// All zero means the vendor cannot do it and the mailbox needs a sign-in.
|
||||
type VendorAuthorization struct {
|
||||
GrantID *uuid.UUID
|
||||
Pending bool
|
||||
// Message says why the vendor could not authorize the domain.
|
||||
Message string
|
||||
}
|
||||
|
||||
// VendorAuthorizer is a VendorSource that can have the vendor authorize this
|
||||
// instance's app on a domain, so its mailboxes connect with no sign-in.
|
||||
type VendorAuthorizer interface {
|
||||
AuthorizeDomain(ctx context.Context, orgID, userID, connectionID uuid.UUID, email, provider string) VendorAuthorization
|
||||
}
|
||||
|
||||
// Deps is everything the service talks to. Asker, Warmup, Auditor,
|
||||
// Publisher, Delegator and Vendors are optional.
|
||||
type Deps struct {
|
||||
|
||||
@@ -0,0 +1,246 @@
|
||||
package vendorconn
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/rs/zerolog/log"
|
||||
|
||||
"github.com/warmbly/warmbly/internal/app/mailboximport"
|
||||
"github.com/warmbly/warmbly/internal/errx"
|
||||
"github.com/warmbly/warmbly/internal/models"
|
||||
"github.com/warmbly/warmbly/internal/pkg/mailvendor"
|
||||
)
|
||||
|
||||
// DomainAuthorization is the import's view of one domain's authorization.
|
||||
type DomainAuthorization = mailboximport.VendorAuthorization
|
||||
|
||||
// authState is one domain's authorization, shared by every row and replica through Redis.
|
||||
type authState struct {
|
||||
RequestID string `json:"request_id,omitempty"`
|
||||
Failed string `json:"failed,omitempty"`
|
||||
Started time.Time `json:"started,omitempty"`
|
||||
Completed time.Time `json:"completed,omitempty"`
|
||||
}
|
||||
|
||||
const (
|
||||
authStartTTL = 2 * time.Minute
|
||||
authPendingTTL = 24 * time.Hour
|
||||
authFailedTTL = time.Hour
|
||||
// authSettle is how long a completed consent may take to reach the provider's token service.
|
||||
authSettle = 15 * time.Minute
|
||||
// authMaxWait bounds one request, so no row waits on a vendor forever.
|
||||
authMaxWait = 2 * time.Hour
|
||||
listTTL = time.Minute
|
||||
)
|
||||
|
||||
func authKey(orgID, connectionID uuid.UUID, provider, domain string) string {
|
||||
return "vendor_app_auth:" + orgID.String() + ":" + connectionID.String() + ":" + provider + ":" + domain
|
||||
}
|
||||
|
||||
// AuthorizeDomain has the vendor authorize this instance's app on a mailbox's
|
||||
// domain and records the grant once it is done; call again while Pending.
|
||||
func (s *Service) AuthorizeDomain(ctx context.Context, orgID, userID, connectionID uuid.UUID, email, provider string) DomainAuthorization {
|
||||
if s.grants == nil || s.cache == nil {
|
||||
return DomainAuthorization{}
|
||||
}
|
||||
at := strings.LastIndex(email, "@")
|
||||
if at < 0 || (provider != models.GrantProviderGoogle && provider != models.GrantProviderMicrosoft) {
|
||||
return DomainAuthorization{}
|
||||
}
|
||||
domain := strings.ToLower(email[at+1:])
|
||||
if g, err := s.grants.GrantFor(ctx, orgID, provider, domain); err == nil && g != nil {
|
||||
return DomainAuthorization{GrantID: &g.ID}
|
||||
}
|
||||
clientID, scopes, roles, ok := s.grants.AppIdentity(provider)
|
||||
if !ok {
|
||||
return DomainAuthorization{}
|
||||
}
|
||||
c, xerr := s.get(ctx, orgID, connectionID)
|
||||
if xerr != nil {
|
||||
return DomainAuthorization{Pending: xerr.Code == errx.Internal}
|
||||
}
|
||||
client, xerr := s.client(ctx, c)
|
||||
if xerr != nil {
|
||||
if xerr.Code == errx.Internal {
|
||||
return DomainAuthorization{Pending: true}
|
||||
}
|
||||
return DomainAuthorization{Message: xerr.Message}
|
||||
}
|
||||
az, ok := client.(mailvendor.AppAuthorizer)
|
||||
if !ok {
|
||||
return DomainAuthorization{}
|
||||
}
|
||||
label := labelOf(c.Vendor)
|
||||
key := authKey(orgID, connectionID, provider, domain)
|
||||
|
||||
var st authState
|
||||
if err := s.cache.GetJSON(ctx, key, &st); err == nil {
|
||||
switch {
|
||||
case st.Failed != "":
|
||||
return DomainAuthorization{Message: st.Failed}
|
||||
case st.RequestID == "":
|
||||
return DomainAuthorization{Pending: true}
|
||||
}
|
||||
status, err := az.AuthorizationStatus(ctx, st.RequestID)
|
||||
switch {
|
||||
case err != nil && transient(err):
|
||||
status = mailvendor.AuthorizationStatus{State: mailvendor.AuthorizationPending}
|
||||
case errors.Is(err, mailvendor.ErrUnauthorized):
|
||||
return DomainAuthorization{Message: s.failed(ctx, c, err).Message}
|
||||
case err != nil:
|
||||
status = mailvendor.AuthorizationStatus{State: mailvendor.AuthorizationFailed, Reason: "the request is gone"}
|
||||
}
|
||||
if status.State == mailvendor.AuthorizationPending && time.Since(st.Started) > authMaxWait {
|
||||
status = mailvendor.AuthorizationStatus{State: mailvendor.AuthorizationFailed, Reason: "it did not finish within two hours"}
|
||||
}
|
||||
switch status.State {
|
||||
case mailvendor.AuthorizationPending:
|
||||
return DomainAuthorization{Pending: true}
|
||||
case mailvendor.AuthorizationFailed:
|
||||
msg := label + " could not authorize Warmbly on " + domain + reasonSuffix(status.Reason) + ". Sign in on each mailbox instead, or retry these rows once it is fixed."
|
||||
_ = s.cache.SetJSON(ctx, key, authState{Failed: msg}, authFailedTTL)
|
||||
return DomainAuthorization{Message: msg}
|
||||
}
|
||||
if st.Completed.IsZero() {
|
||||
st.Completed = time.Now()
|
||||
_ = s.cache.SetJSON(ctx, key, st, authPendingTTL)
|
||||
}
|
||||
return s.recordGrant(ctx, c, client, key, orgID, userID, provider, domain, time.Since(st.Completed) < authSettle)
|
||||
}
|
||||
|
||||
started, err := s.cache.SetNX(ctx, key, "{}", authStartTTL).Result()
|
||||
if err != nil {
|
||||
return DomainAuthorization{}
|
||||
}
|
||||
if !started {
|
||||
return DomainAuthorization{Pending: true}
|
||||
}
|
||||
list, err := s.listed(ctx, c, client)
|
||||
if err != nil {
|
||||
_ = s.cache.Del(ctx, key)
|
||||
switch {
|
||||
case transient(err):
|
||||
return DomainAuthorization{Pending: true}
|
||||
case errors.Is(err, mailvendor.ErrUnauthorized):
|
||||
return DomainAuthorization{Message: s.failed(ctx, c, err).Message}
|
||||
}
|
||||
return DomainAuthorization{}
|
||||
}
|
||||
// The admin mailbox names the workspace InboxKit resolves the domain's admin in.
|
||||
var on *mailvendor.Mailbox
|
||||
admin := false
|
||||
for i := range list {
|
||||
if strings.EqualFold(list[i].Domain, domain) || strings.HasSuffix(strings.ToLower(list[i].Email), "@"+domain) {
|
||||
if on == nil || (list[i].Admin && !admin) {
|
||||
on = &list[i]
|
||||
}
|
||||
admin = admin || list[i].Admin
|
||||
}
|
||||
}
|
||||
// A Google grant is verified as the domain's administrator, so without one there is nothing to verify.
|
||||
if on == nil || (provider == models.GrantProviderGoogle && !admin) {
|
||||
_ = s.cache.Del(ctx, key)
|
||||
return DomainAuthorization{}
|
||||
}
|
||||
reqID, err := az.AuthorizeApp(ctx, mailvendor.AppAuthorization{
|
||||
Domain: domain, Mailbox: *on, Provider: provider, ClientID: clientID, Scopes: scopes, AppRoles: roles,
|
||||
})
|
||||
if err != nil {
|
||||
if transient(err) {
|
||||
_ = s.cache.Del(ctx, key)
|
||||
return DomainAuthorization{Pending: true}
|
||||
}
|
||||
if errors.Is(err, mailvendor.ErrUnauthorized) {
|
||||
_ = s.failed(ctx, c, err)
|
||||
}
|
||||
log.Info().Str("vendor", c.Vendor).Str("domain", domain).Err(err).Msg("vendor connection: app authorization refused")
|
||||
msg := label + " could not authorize Warmbly on " + domain + ". It needs an admin mailbox on the domain. Sign in on each mailbox instead."
|
||||
_ = s.cache.SetJSON(ctx, key, authState{Failed: msg}, authFailedTTL)
|
||||
return DomainAuthorization{Message: msg}
|
||||
}
|
||||
_ = s.cache.SetJSON(ctx, key, authState{RequestID: reqID, Started: time.Now()}, authPendingTTL)
|
||||
return DomainAuthorization{Pending: true}
|
||||
}
|
||||
|
||||
// recordGrant turns a completed vendor authorization into the workspace's grant,
|
||||
// covering only domains this vendor account holds.
|
||||
// settling keeps a refusal pending while a fresh consent propagates.
|
||||
func (s *Service) recordGrant(ctx context.Context, c *models.VendorConnection, client mailvendor.Client, key string, orgID, userID uuid.UUID, provider, domain string, settling bool) DomainAuthorization {
|
||||
list, err := s.listed(ctx, c, client)
|
||||
if err != nil {
|
||||
return DomainAuthorization{Pending: true}
|
||||
}
|
||||
owned := map[string]bool{}
|
||||
admin := ""
|
||||
for _, m := range list {
|
||||
d := strings.ToLower(m.Domain)
|
||||
if at := strings.LastIndex(m.Email, "@"); d == "" && at > 0 {
|
||||
d = strings.ToLower(m.Email[at+1:])
|
||||
}
|
||||
if d != "" {
|
||||
owned[d] = true
|
||||
}
|
||||
if m.Admin && d == domain && admin == "" {
|
||||
admin = m.Email
|
||||
}
|
||||
}
|
||||
domains := make([]string, 0, len(owned))
|
||||
for d := range owned {
|
||||
domains = append(domains, d)
|
||||
}
|
||||
g, xerr := s.grants.GrantFromVendor(ctx, orgID, userID, provider, domain, admin, domains)
|
||||
if xerr != nil {
|
||||
if settling || xerr.Code == errx.Internal || xerr.ResponseCode() == "mailbox_grant_unavailable" {
|
||||
return DomainAuthorization{Pending: true}
|
||||
}
|
||||
msg := labelOf(c.Vendor) + " authorized Warmbly on " + domain + ", but the grant did not verify: " + xerr.Message
|
||||
_ = s.cache.SetJSON(ctx, key, authState{Failed: msg}, authFailedTTL)
|
||||
return DomainAuthorization{Message: msg}
|
||||
}
|
||||
_ = s.cache.Del(ctx, key)
|
||||
return DomainAuthorization{GrantID: &g.ID}
|
||||
}
|
||||
|
||||
// listed is the connection's mailbox list, shared a minute so every domain of an import reads it once.
|
||||
func (s *Service) listed(ctx context.Context, c *models.VendorConnection, client mailvendor.Client) ([]mailvendor.Mailbox, error) {
|
||||
if v, ok := s.lists.Load(c.ID); ok {
|
||||
if l := v.(cachedList); time.Since(l.at) < listTTL {
|
||||
return l.boxes, nil
|
||||
}
|
||||
}
|
||||
boxes, err := client.List(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
s.lists.Store(c.ID, cachedList{boxes: boxes, at: time.Now()})
|
||||
return boxes, nil
|
||||
}
|
||||
|
||||
type cachedList struct {
|
||||
boxes []mailvendor.Mailbox
|
||||
at time.Time
|
||||
}
|
||||
|
||||
// transient is a vendor failure worth trying again: nothing was decided.
|
||||
func transient(err error) bool {
|
||||
if errors.Is(err, mailvendor.ErrRateLimited) || errors.Is(err, context.DeadlineExceeded) || errors.Is(err, context.Canceled) {
|
||||
return true
|
||||
}
|
||||
var ve *mailvendor.Error
|
||||
return errors.As(err, &ve) && (ve.Status == 0 || ve.Status >= 500)
|
||||
}
|
||||
|
||||
func reasonSuffix(reason string) string {
|
||||
reason = strings.TrimSpace(reason)
|
||||
if reason == "" {
|
||||
return ""
|
||||
}
|
||||
if r := []rune(reason); len(r) > 200 {
|
||||
reason = string(r[:200])
|
||||
}
|
||||
return " (" + strings.TrimRight(reason, ".") + ")"
|
||||
}
|
||||
@@ -49,6 +49,14 @@ type Reconnector interface {
|
||||
UpdateSMTPIMAPCredentials(ctx context.Context, orgID *uuid.UUID, accountID uuid.UUID, creds *models.SmtpImap) (*models.Email, *errx.Error)
|
||||
}
|
||||
|
||||
// Grants is the delegation side: the instance's app identity, and recording a
|
||||
// grant once a vendor has authorized that app on a domain.
|
||||
type Grants interface {
|
||||
AppIdentity(provider string) (clientID string, scopes, appRoles []string, ok bool)
|
||||
GrantFor(ctx context.Context, orgID uuid.UUID, provider, domain string) (*models.DomainGrant, error)
|
||||
GrantFromVendor(ctx context.Context, orgID, userID uuid.UUID, provider, domain, admin string, owned []string) (*models.DomainGrant, *errx.Error)
|
||||
}
|
||||
|
||||
// Importer starts a background import of picked mailboxes.
|
||||
type Importer interface {
|
||||
CreateFromList(ctx context.Context, in mailboximport.ListInput) (*models.MailboxImport, *errx.Error)
|
||||
@@ -60,6 +68,7 @@ type Service struct {
|
||||
mailboxes Mailboxes
|
||||
reconnect Reconnector
|
||||
importer Importer
|
||||
grants Grants
|
||||
cache *cache.Cache
|
||||
newClient func(vendor string, fields map[string]string) (mailvendor.Client, error)
|
||||
|
||||
@@ -68,6 +77,8 @@ type Service struct {
|
||||
clients sync.Map // connection id + credential hash -> *cachedClient
|
||||
// domainCache keeps each connection's domain list a few minutes.
|
||||
domainCache domainCache
|
||||
// lists keeps each connection's mailbox list a minute, for app authorization.
|
||||
lists sync.Map // connection id -> cachedList
|
||||
}
|
||||
|
||||
type cachedClient struct {
|
||||
@@ -83,13 +94,14 @@ type Deps struct {
|
||||
Mailboxes Mailboxes
|
||||
Reconnect Reconnector
|
||||
Importer Importer
|
||||
Grants Grants
|
||||
Cache *cache.Cache
|
||||
// NewClient overrides the vendor client (tests).
|
||||
NewClient func(vendor string, fields map[string]string) (mailvendor.Client, error)
|
||||
}
|
||||
|
||||
func NewService(d Deps) *Service {
|
||||
s := &Service{repo: d.Repo, cipher: d.Cipher, mailboxes: d.Mailboxes, reconnect: d.Reconnect, importer: d.Importer, cache: d.Cache, newClient: d.NewClient}
|
||||
s := &Service{repo: d.Repo, cipher: d.Cipher, mailboxes: d.Mailboxes, reconnect: d.Reconnect, importer: d.Importer, grants: d.Grants, cache: d.Cache, newClient: d.NewClient}
|
||||
if s.newClient == nil {
|
||||
s.newClient = func(vendor string, fields map[string]string) (mailvendor.Client, error) {
|
||||
return mailvendor.New(vendor, fields)
|
||||
@@ -338,6 +350,7 @@ func (s *Service) client(ctx context.Context, c *models.VendorConnection) (mailv
|
||||
// forget drops a connection's cached clients and domain list, so a changed or removed key is never used again.
|
||||
func (s *Service) forget(id uuid.UUID) {
|
||||
s.forgetDomains(id)
|
||||
s.lists.Delete(id)
|
||||
prefix := id.String() + ":"
|
||||
s.clients.Range(func(k, _ any) bool {
|
||||
if strings.HasPrefix(k.(string), prefix) {
|
||||
|
||||
@@ -0,0 +1,131 @@
|
||||
package worker
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net"
|
||||
"strconv"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
|
||||
"github.com/warmbly/warmbly/internal/app/worker/mailmanager"
|
||||
"github.com/warmbly/warmbly/internal/models"
|
||||
)
|
||||
|
||||
type capturedEvents struct {
|
||||
mu sync.Mutex
|
||||
events []models.JobEventType
|
||||
bodies []any
|
||||
}
|
||||
|
||||
func (c *capturedEvents) on(t models.JobEventType, _ string, body any) error {
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
c.events, c.bodies = append(c.events, t), append(c.bodies, body)
|
||||
return nil
|
||||
}
|
||||
|
||||
func newLoadingWorker(c *capturedEvents) *WorkerService {
|
||||
return &WorkerService{mailManager: mailmanager.NewMailManager(c.on, nil, nil, nil, nil, nil, nil)}
|
||||
}
|
||||
|
||||
// A mail server that accepts and never speaks holds a load for minutes; it must not hold the command queue.
|
||||
func TestAddEmailDoesNotBlockOnASilentServer(t *testing.T) {
|
||||
ln, err := net.Listen("tcp", "127.0.0.1:0")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer ln.Close()
|
||||
var accepted atomic.Int32
|
||||
go func() {
|
||||
for {
|
||||
conn, err := ln.Accept()
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
accepted.Add(1)
|
||||
defer conn.Close()
|
||||
}
|
||||
}()
|
||||
port, _ := strconv.Atoi(strconv.Itoa(ln.Addr().(*net.TCPAddr).Port))
|
||||
|
||||
w := newLoadingWorker(&capturedEvents{})
|
||||
id := uuid.New()
|
||||
e := &models.AddWorkerEmail{
|
||||
ID: id, UserID: uuid.New(), Type: models.InboxProviderSMTPIMAP, ImapSync: true, Email: "a@example.test",
|
||||
SmtpImap: &models.AddWorkerEmailSmtpImapData{Credentials: &models.SmtpImap{
|
||||
IMAP: &models.Service{Host: "127.0.0.1", Port: port, Username: "a", Password: "b", Security: models.MailSecurityTLS},
|
||||
SMTP: &models.Service{Host: "127.0.0.1", Port: port, Username: "a", Password: "b", Security: models.MailSecurityTLS},
|
||||
}},
|
||||
}
|
||||
start := time.Now()
|
||||
for i := 0; i < 3; i++ {
|
||||
if err := w.HandleAddEmail(context.Background(), e); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
if took := time.Since(start); took > time.Second {
|
||||
t.Fatalf("three ADD_EMAILs held the queue for %v", took)
|
||||
}
|
||||
time.Sleep(200 * time.Millisecond)
|
||||
if n := accepted.Load(); n != 1 {
|
||||
t.Fatalf("a republished ADD_EMAIL dialed again: %d connections", n)
|
||||
}
|
||||
|
||||
// A command for the mailbox waits for the load, within its own deadline.
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 100*time.Millisecond)
|
||||
defer cancel()
|
||||
if _, ok := w.loadedMailbox(ctx, id); ok {
|
||||
t.Fatal("a mailbox still dialing was handed out")
|
||||
}
|
||||
if err := w.HandleRemoveEmail(context.Background(), &models.RemoveWorkerEmail{EmailID: id.String()}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if v, ok := w.loads.Load(id); !ok || !v.(*mailboxLoad).removed.Load() {
|
||||
t.Fatal("a removal during the load was not recorded, so the mailbox would load after it was deleted")
|
||||
}
|
||||
}
|
||||
|
||||
func TestAddEmailReportsAMailboxThatCannotLoad(t *testing.T) {
|
||||
c := &capturedEvents{}
|
||||
w := newLoadingWorker(c)
|
||||
id := uuid.New()
|
||||
if err := w.HandleAddEmail(context.Background(), &models.AddWorkerEmail{ID: id, UserID: uuid.New(), Type: models.InboxProviderSMTPIMAP}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
|
||||
defer cancel()
|
||||
if _, ok := w.loadedMailbox(ctx, id); ok {
|
||||
t.Fatal("a mailbox without credentials loaded")
|
||||
}
|
||||
c.mu.Lock()
|
||||
defer c.mu.Unlock()
|
||||
if len(c.events) != 1 || c.events[0] != models.JobEventTypeEmailAuthError {
|
||||
t.Fatalf("events = %v, want one auth error so the mailbox shows it and stops", c.events)
|
||||
}
|
||||
if ev, ok := c.bodies[0].(models.EmailErrorEvent); !ok || ev.EmailAccountID != id.String() {
|
||||
t.Fatalf("event body = %+v", c.bodies[0])
|
||||
}
|
||||
}
|
||||
|
||||
func TestAddEmailLoadsAndCommandsFindIt(t *testing.T) {
|
||||
w := newLoadingWorker(&capturedEvents{})
|
||||
id := uuid.New()
|
||||
e := &models.AddWorkerEmail{
|
||||
ID: id, UserID: uuid.New(), Type: models.InboxProviderSMTPIMAP, Email: "a@example.test",
|
||||
SmtpImap: &models.AddWorkerEmailSmtpImapData{Credentials: &models.SmtpImap{
|
||||
IMAP: &models.Service{Host: "imap.example.test", Port: 993}, SMTP: &models.Service{Host: "smtp.example.test", Port: 587},
|
||||
}},
|
||||
}
|
||||
if err := w.HandleAddEmail(context.Background(), e); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
|
||||
defer cancel()
|
||||
if _, ok := w.loadedMailbox(ctx, id); !ok {
|
||||
t.Fatal("a command right behind the ADD_EMAIL did not find the mailbox")
|
||||
}
|
||||
}
|
||||
@@ -62,10 +62,9 @@ type WorkerAssignmentService interface {
|
||||
// MigrateEmailsFromWorker drains every mailbox off a worker.
|
||||
MigrateEmailsFromWorker(ctx context.Context, workerID uuid.UUID) error
|
||||
|
||||
// SelectValidationWorker returns any live worker to run a one-shot
|
||||
// credential handshake on. Nothing is placed, so no scoring applies: the
|
||||
// worker only dials the mailbox once and reports back.
|
||||
SelectValidationWorker(ctx context.Context) (*models.Worker, error)
|
||||
// ValidationWorkers orders the healthy live workers for a credential check:
|
||||
// prefer, then where placement would put orgID's new mailbox, then by load.
|
||||
ValidationWorkers(ctx context.Context, orgID uuid.UUID, prefer *uuid.UUID) ([]models.Worker, error)
|
||||
}
|
||||
|
||||
// SetOrganizationWarmupPool idempotently applies a subscription tier to existing mailboxes.
|
||||
@@ -506,15 +505,39 @@ func (s *workerAssignmentService) MigrateEmailsFromWorker(ctx context.Context, w
|
||||
return nil
|
||||
}
|
||||
|
||||
// SelectValidationWorker returns the least loaded live worker. Used by the
|
||||
// connect and reconnect flows to test credentials before anything is stored.
|
||||
func (s *workerAssignmentService) SelectValidationWorker(ctx context.Context) (*models.Worker, error) {
|
||||
workers, err := s.workerRepo.ListPlaceableWorkers(ctx)
|
||||
// ValidationWorkers puts the address the mailbox will sign in from first, so the provider sees one IP.
|
||||
func (s *workerAssignmentService) ValidationWorkers(ctx context.Context, orgID uuid.UUID, prefer *uuid.UUID) ([]models.Worker, error) {
|
||||
live, err := s.workerRepo.ListPlaceableWorkers(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if len(workers) == 0 {
|
||||
healthy := make([]models.Worker, 0, len(live))
|
||||
for _, w := range live {
|
||||
switch w.HealthState {
|
||||
case models.WorkerHealthHealthy, models.WorkerHealthWatch:
|
||||
healthy = append(healthy, w)
|
||||
}
|
||||
}
|
||||
if len(healthy) == 0 {
|
||||
return nil, ErrNoAvailableWorkers
|
||||
}
|
||||
return &workers[0], nil
|
||||
first := func(id uuid.UUID) {
|
||||
for i := range healthy {
|
||||
if healthy[i].ID == id {
|
||||
w := healthy[i]
|
||||
healthy = append(healthy[:i], healthy[i+1:]...)
|
||||
healthy = append([]models.Worker{w}, healthy...)
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
if orgID != uuid.Nil {
|
||||
if res, err := s.SelectWorkerFor(ctx, PlacementLookup{OrgID: orgID}); err == nil && res != nil && res.Worker != nil {
|
||||
first(res.Worker.ID)
|
||||
}
|
||||
}
|
||||
if prefer != nil {
|
||||
first(*prefer)
|
||||
}
|
||||
return healthy, nil
|
||||
}
|
||||
|
||||
@@ -2,11 +2,27 @@ package worker
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/rs/zerolog/log"
|
||||
"github.com/warmbly/warmbly/internal/app/worker/wmail"
|
||||
"github.com/warmbly/warmbly/internal/errx"
|
||||
"github.com/warmbly/warmbly/internal/models"
|
||||
"github.com/warmbly/warmbly/internal/observability/errs"
|
||||
)
|
||||
|
||||
// mailboxLoad is one mailbox being loaded; removed is set when a REMOVE_EMAIL arrives meanwhile.
|
||||
type mailboxLoad struct {
|
||||
removed atomic.Bool
|
||||
done chan struct{}
|
||||
}
|
||||
|
||||
// HandleAddEmail loads a mailbox off the bus loop. Loading dials the mail
|
||||
// server, and one that is slow or refuses would otherwise hold every command
|
||||
// queued behind it, credential checks and sends included.
|
||||
func (w *WorkerService) HandleAddEmail(ctx context.Context, e *models.AddWorkerEmail) error {
|
||||
if e == nil {
|
||||
return nil
|
||||
@@ -21,24 +37,101 @@ func (w *WorkerService) HandleAddEmail(ctx context.Context, e *models.AddWorkerE
|
||||
return nil
|
||||
}
|
||||
|
||||
if err := w.mailManager.AddWMail(ctx, e); err != nil {
|
||||
log.Error().Err(err).Str("email_id", e.ID.String()).Msg("failed to add email account to worker")
|
||||
return err
|
||||
load := &mailboxLoad{done: make(chan struct{})}
|
||||
if _, busy := w.loads.LoadOrStore(e.ID, load); busy {
|
||||
return nil
|
||||
}
|
||||
|
||||
// Start the periodic mail sync worker. Uses the WMail's own context which is
|
||||
// cancelled when the account is removed or terminates, so we don't leak goroutines.
|
||||
mail := w.mailManager.Get(e.ID)
|
||||
if mail != nil {
|
||||
// Gmail (history) and Outlook/Graph (delta) always sync. Only generic
|
||||
// SMTP/IMAP mailboxes are opt-in via ImapSync.
|
||||
if e.Type == models.InboxProviderSMTPIMAP && !e.ImapSync {
|
||||
log.Info().Str("email_id", e.ID.String()).Str("email", e.Email).Msg("email account added (no sync)")
|
||||
return nil
|
||||
}
|
||||
go mail.StartSyncWorker(mail.Ctx)
|
||||
}
|
||||
|
||||
log.Info().Str("email_id", e.ID.String()).Str("email", e.Email).Msg("email account added to worker")
|
||||
go func() {
|
||||
defer close(load.done)
|
||||
defer w.loads.Delete(e.ID)
|
||||
defer func() {
|
||||
if r := recover(); r != nil {
|
||||
errs.Recover(r)
|
||||
}
|
||||
}()
|
||||
w.loadMailbox(context.WithoutCancel(ctx), e, load)
|
||||
}()
|
||||
return nil
|
||||
}
|
||||
|
||||
// loadedMailbox returns a mailbox, waiting on a load of it still in flight so a
|
||||
// command queued behind its ADD_EMAIL finds it as it did when loads were inline.
|
||||
func (w *WorkerService) loadedMailbox(ctx context.Context, id uuid.UUID) (*wmail.WMail, bool) {
|
||||
if v, ok := w.loads.Load(id); ok {
|
||||
select {
|
||||
case <-v.(*mailboxLoad).done:
|
||||
case <-ctx.Done():
|
||||
}
|
||||
}
|
||||
mail := w.mailManager.Get(id)
|
||||
return mail, mail != nil
|
||||
}
|
||||
|
||||
// loadMailbox connects one mailbox and starts its sync, or reports why it could not.
|
||||
func (w *WorkerService) loadMailbox(ctx context.Context, e *models.AddWorkerEmail, load *mailboxLoad) {
|
||||
if err := w.mailManager.AddWMail(ctx, e); err != nil {
|
||||
log.Error().Err(err).Str("email_id", e.ID.String()).Msg("failed to add email account to worker")
|
||||
w.reportLoadFailure(e, err)
|
||||
return
|
||||
}
|
||||
mail := w.mailManager.Get(e.ID)
|
||||
if mail == nil {
|
||||
return
|
||||
}
|
||||
if load.removed.Load() {
|
||||
mail.Discard()
|
||||
w.mailManager.Terminate(e.ID)
|
||||
return
|
||||
}
|
||||
|
||||
// Gmail (history) and Outlook/Graph (delta) always sync. Only generic
|
||||
// SMTP/IMAP mailboxes are opt-in via ImapSync. The sync worker runs on the
|
||||
// WMail's own context, cancelled when the account is removed or terminates.
|
||||
if e.Type == models.InboxProviderSMTPIMAP && !e.ImapSync {
|
||||
log.Info().Str("email_id", e.ID.String()).Str("email", e.Email).Msg("email account added (no sync)")
|
||||
return
|
||||
}
|
||||
go mail.StartSyncWorker(mail.Ctx)
|
||||
log.Info().Str("email_id", e.ID.String()).Str("email", e.Email).Msg("email account added to worker")
|
||||
}
|
||||
|
||||
// reportLoadFailure raises the same account event a sync failure would, so a
|
||||
// mailbox that cannot sign in shows the error and stops instead of looking connected.
|
||||
func (w *WorkerService) reportLoadFailure(e *models.AddWorkerEmail, err error) {
|
||||
var mailErr *errx.MailError
|
||||
if !errors.As(err, &mailErr) {
|
||||
return
|
||||
}
|
||||
eventType := wmail.DetermineErrorEventType(mailErr)
|
||||
switch eventType {
|
||||
case models.JobEventTypeEmailAuthError, models.JobEventTypeEmailDisabled,
|
||||
models.JobEventTypeEmailRateLimited, models.JobEventTypeEmailServerError:
|
||||
default:
|
||||
return
|
||||
}
|
||||
info := mailErr.GetUserErrorInfo()
|
||||
event := models.EmailErrorEvent{
|
||||
EmailAccountID: e.ID.String(),
|
||||
UserID: e.UserID.String(),
|
||||
ErrorCode: string(mailErr.Code),
|
||||
ErrorType: string(mailErr.Type),
|
||||
ResolveMethod: string(mailErr.ResolveMethod),
|
||||
Message: mailErr.Message,
|
||||
UserVisible: mailErr.IsUserVisible(),
|
||||
UserTitle: info.Title,
|
||||
UserMessage: info.Message,
|
||||
ActionRequired: info.ActionRequired,
|
||||
Timestamp: time.Now().Unix(),
|
||||
}
|
||||
if perr := w.mailManager.OnEvent(eventType, e.ID.String(), event); perr != nil {
|
||||
log.Error().Err(perr).Str("email_id", e.ID.String()).Msg("failed to report a mailbox that did not load")
|
||||
}
|
||||
}
|
||||
|
||||
// dropMailbox removes a loaded mailbox and stops its work.
|
||||
func (w *WorkerService) dropMailbox(id uuid.UUID, mail *wmail.WMail) {
|
||||
if mail.Cancel != nil {
|
||||
mail.Cancel()
|
||||
}
|
||||
w.mailManager.Terminate(id)
|
||||
}
|
||||
|
||||
@@ -48,6 +48,33 @@ func (w *WorkerService) HandleEmailValidation(parent context.Context, data model
|
||||
return nil
|
||||
}
|
||||
|
||||
// ListenValidations takes credential checks from this worker's Redis request
|
||||
// channel until ctx ends. go-redis resubscribes by itself after a dropped link.
|
||||
func (w *WorkerService) ListenValidations(ctx context.Context) {
|
||||
id, err := uuid.Parse(w.ID)
|
||||
if err != nil || w.Cache == nil {
|
||||
return
|
||||
}
|
||||
sub := w.Cache.Subscribe(ctx, models.EmailValidationRequestChannel(id))
|
||||
defer sub.Close()
|
||||
msgs := sub.Channel()
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case msg, ok := <-msgs:
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
var data models.EventWorkerEmailValidation
|
||||
if err := json.Unmarshal([]byte(msg.Payload), &data); err != nil || data.ProcessID == uuid.Nil {
|
||||
continue
|
||||
}
|
||||
_ = w.HandleEmailValidation(ctx, data)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// runEmailValidation unseals the credentials, probes both legs and always
|
||||
// answers: a check the worker could not run is reported as such, because a
|
||||
// silent worker reads at the backend as a mail server that never replied.
|
||||
@@ -123,7 +150,7 @@ func (w *WorkerService) runEmailValidation(parent context.Context, data models.E
|
||||
func (w *WorkerService) publishValidationVerdict(parent context.Context, processID uuid.UUID, verdict models.EmailValidationVerdict) {
|
||||
replyCtx, replyCancel := context.WithTimeout(context.WithoutCancel(parent), replyBudget)
|
||||
defer replyCancel()
|
||||
channel := "email_validation:" + processID.String()
|
||||
channel := models.EmailValidationReplyChannel(processID)
|
||||
// The verdict goes first and the legacy digit after it: a backend that
|
||||
// reads verdicts returns on the first message, and one that predates them
|
||||
// skips what it cannot parse and takes the digit, so a fleet mid-update
|
||||
|
||||
@@ -46,9 +46,7 @@ func (w *WorkerService) HandleMailboxIdentity(ctx context.Context, data models.E
|
||||
return nil
|
||||
}
|
||||
|
||||
w.mailManager.RLock()
|
||||
mail, exists := w.mailManager.Emails[data.EmailID]
|
||||
w.mailManager.RUnlock()
|
||||
mail, exists := w.loadedMailbox(ctx, data.EmailID)
|
||||
if !exists {
|
||||
// Nothing to answer with: this worker is not running the mailbox. The
|
||||
// backend turns the silence-shaped answer into a message that tells
|
||||
|
||||
@@ -20,9 +20,7 @@ func (w *WorkerService) HandleMessageSeen(ctx context.Context, action models.Mes
|
||||
return nil
|
||||
}
|
||||
|
||||
w.mailManager.RLock()
|
||||
mail, exists := w.mailManager.Emails[action.EmailID]
|
||||
w.mailManager.RUnlock()
|
||||
mail, exists := w.loadedMailbox(ctx, action.EmailID)
|
||||
if !exists {
|
||||
// The mailbox is not loaded here (a restart, or it moved worker mid
|
||||
// flight). The read state is already correct in Warmbly; only the
|
||||
|
||||
@@ -19,17 +19,17 @@ func (w *WorkerService) HandleRemoveEmail(ctx context.Context, e *models.RemoveW
|
||||
return nil
|
||||
}
|
||||
|
||||
// A load still dialing drops the mailbox when it finishes.
|
||||
if v, ok := w.loads.Load(id); ok {
|
||||
v.(*mailboxLoad).removed.Store(true)
|
||||
}
|
||||
|
||||
mail := w.mailManager.Get(id)
|
||||
if mail == nil {
|
||||
// Already gone - idempotent
|
||||
return nil
|
||||
}
|
||||
|
||||
// Cancel any in-flight syncs and remove from manager
|
||||
if mail.Cancel != nil {
|
||||
mail.Cancel()
|
||||
}
|
||||
w.mailManager.Terminate(id)
|
||||
w.dropMailbox(id, mail)
|
||||
|
||||
log.Info().Str("email_id", e.EmailID).Msg("email account removed from worker")
|
||||
return nil
|
||||
|
||||
@@ -25,9 +25,7 @@ func (w *WorkerService) HandleSendEmail(ctx context.Context, sendEmail models.Se
|
||||
Msg("Processing send email event")
|
||||
|
||||
// Get the email account from MailManager
|
||||
w.mailManager.RLock()
|
||||
mail, exists := w.mailManager.Emails[sendEmail.EmailID]
|
||||
w.mailManager.RUnlock()
|
||||
mail, exists := w.loadedMailbox(ctx, sendEmail.EmailID)
|
||||
|
||||
if !exists {
|
||||
// The mailbox is not loaded here: it is still being added (its
|
||||
|
||||
@@ -38,9 +38,7 @@ func (w *WorkerService) HandleWarmupAction(ctx context.Context, action models.Wa
|
||||
return nil
|
||||
}
|
||||
|
||||
w.mailManager.RLock()
|
||||
mail, exists := w.mailManager.Emails[action.EmailID]
|
||||
w.mailManager.RUnlock()
|
||||
mail, exists := w.loadedMailbox(ctx, action.EmailID)
|
||||
var err error
|
||||
switch {
|
||||
case !exists:
|
||||
|
||||
@@ -15,9 +15,6 @@ func (m *MailManager) AddWMail(
|
||||
ctx context.Context,
|
||||
data *models.AddWorkerEmail,
|
||||
) error {
|
||||
m.Lock()
|
||||
defer m.Unlock()
|
||||
|
||||
// Cfg is avro-excluded from the payload, so rebuild it from the worker's
|
||||
// local oauth config for token refresh (no-op for smtp_imap).
|
||||
data.Cfg = m.cfgFor(data.Type)
|
||||
@@ -47,6 +44,14 @@ func (m *MailManager) AddWMail(
|
||||
return err
|
||||
}
|
||||
|
||||
// Built before the lock: NewWMail dials the server, and holding the lock
|
||||
// through that stalled every send and lookup on this worker.
|
||||
m.Lock()
|
||||
defer m.Unlock()
|
||||
if _, loaded := m.Emails[data.ID]; loaded {
|
||||
newMail.Discard()
|
||||
return nil
|
||||
}
|
||||
m.Emails[data.ID] = newMail
|
||||
|
||||
return nil
|
||||
|
||||
@@ -39,6 +39,8 @@ type WorkerService struct {
|
||||
TokenBroker repository.BrokeredTokenClient
|
||||
|
||||
mailManager *mailmanager.MailManager
|
||||
// loads holds the mailboxes being loaded, so a republished ADD_EMAIL never dials twice.
|
||||
loads sync.Map // uuid.UUID -> *mailboxLoad
|
||||
|
||||
// HealthCounters tracks the per-window send-side telemetry the worker
|
||||
// reports via JobEventTypeWorkerHealth. Lazily initialised by
|
||||
|
||||
@@ -0,0 +1,96 @@
|
||||
package worker
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"testing"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/warmbly/warmbly/internal/models"
|
||||
"github.com/warmbly/warmbly/internal/repository"
|
||||
)
|
||||
|
||||
type validationRepo struct {
|
||||
repository.WorkerRepository
|
||||
live []models.Worker
|
||||
placed *uuid.UUID
|
||||
}
|
||||
|
||||
func (r *validationRepo) ListPlaceableWorkers(context.Context) ([]models.Worker, error) {
|
||||
return append([]models.Worker(nil), r.live...), nil
|
||||
}
|
||||
func (r *validationRepo) GetEmailAccountPlacementHint(context.Context, uuid.UUID) (*repository.EmailAccountPlacementHint, error) {
|
||||
return nil, errors.New("no mailbox yet")
|
||||
}
|
||||
func (r *validationRepo) CountOrgMailboxes(context.Context, uuid.UUID) (int, error) { return 0, nil }
|
||||
func (r *validationRepo) ListPlacementCandidates(context.Context, uuid.UUID, string, []models.WorkerHealthState) ([]repository.PlacementCandidateRow, error) {
|
||||
if r.placed == nil {
|
||||
return nil, nil
|
||||
}
|
||||
return []repository.PlacementCandidateRow{{WorkerCapacityRowDB: repository.WorkerCapacityRowDB{
|
||||
WorkerID: *r.placed, HealthState: models.WorkerHealthHealthy, BaseCapacity: 16, HealthMultiplier: 1, AgeMultiplier: 1,
|
||||
}}}, nil
|
||||
}
|
||||
func (r *validationRepo) GetByID(_ context.Context, id uuid.UUID) (*models.Worker, error) {
|
||||
for _, w := range r.live {
|
||||
if w.ID == id {
|
||||
return &w, nil
|
||||
}
|
||||
}
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
func ids(ws []models.Worker) []uuid.UUID {
|
||||
out := make([]uuid.UUID, len(ws))
|
||||
for i, w := range ws {
|
||||
out[i] = w.ID
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
func TestValidationWorkersOrderAndHealth(t *testing.T) {
|
||||
a, b, c, blocked := uuid.New(), uuid.New(), uuid.New(), uuid.New()
|
||||
repo := &validationRepo{live: []models.Worker{
|
||||
{ID: a, HealthState: models.WorkerHealthHealthy},
|
||||
{ID: blocked, HealthState: models.WorkerHealthBlocked},
|
||||
{ID: b, HealthState: models.WorkerHealthWatch},
|
||||
{ID: c, HealthState: models.WorkerHealthHealthy},
|
||||
}, placed: &c}
|
||||
svc := &workerAssignmentService{workerRepo: repo}
|
||||
org := uuid.New()
|
||||
|
||||
got, err := svc.ValidationWorkers(context.Background(), org, nil)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if want := []uuid.UUID{c, a, b}; !equalIDs(ids(got), want) {
|
||||
t.Fatalf("new mailbox order = %v, want placement first then by load %v", ids(got), want)
|
||||
}
|
||||
|
||||
got, _ = svc.ValidationWorkers(context.Background(), org, &b)
|
||||
if want := []uuid.UUID{b, c, a}; !equalIDs(ids(got), want) {
|
||||
t.Fatalf("reconnect order = %v, want the mailbox's own worker first %v", ids(got), want)
|
||||
}
|
||||
|
||||
got, _ = svc.ValidationWorkers(context.Background(), org, &blocked)
|
||||
if want := []uuid.UUID{c, a, b}; !equalIDs(ids(got), want) {
|
||||
t.Fatalf("an unhealthy own worker was offered: %v", ids(got))
|
||||
}
|
||||
|
||||
repo.live = []models.Worker{{ID: blocked, HealthState: models.WorkerHealthBlocked}}
|
||||
if _, err := svc.ValidationWorkers(context.Background(), org, nil); !errors.Is(err, ErrNoAvailableWorkers) {
|
||||
t.Fatalf("no healthy worker: err = %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func equalIDs(a, b []uuid.UUID) bool {
|
||||
if len(a) != len(b) {
|
||||
return false
|
||||
}
|
||||
for i := range a {
|
||||
if a[i] != b[i] {
|
||||
return false
|
||||
}
|
||||
}
|
||||
return true
|
||||
}
|
||||
@@ -321,3 +321,16 @@ func (w *WMail) ApplySyncPolicy(data *models.AddWorkerEmailSyncData) {
|
||||
}
|
||||
w.gov.SetPolicy(data.Policy)
|
||||
}
|
||||
|
||||
// Discard releases a WMail that was built but never put to work.
|
||||
func (w *WMail) Discard() {
|
||||
if w.Cancel != nil {
|
||||
w.Cancel()
|
||||
}
|
||||
if w.SmtpImapData == nil || w.SmtpImapData.ImapClient == nil {
|
||||
return
|
||||
}
|
||||
if c, ok := w.SmtpImapData.ImapClient.(interface{ Close() error }); ok {
|
||||
_ = c.Close()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -102,8 +102,8 @@ func (b *KafkaBus) Publish(ctx context.Context, topic, key string, payload []byt
|
||||
|
||||
// Subscribe creates a fresh consumer in the given group, subscribes to all
|
||||
// topics, and blocks reading messages until ctx is cancelled or a fatal error
|
||||
// occurs. Handler errors are logged but do not abort the loop; the message is
|
||||
// still committed to match the existing kafka.Consumer.Consume behaviour.
|
||||
// occurs. Handler errors are logged but do not abort the loop; the message's
|
||||
// offset is still stored, to match the existing kafka.Consumer.Consume behaviour.
|
||||
func (b *KafkaBus) Subscribe(ctx context.Context, topics []string, group string, handler Handler) error {
|
||||
if len(topics) == 0 {
|
||||
return errors.New("eventbus kafka: at least one topic required")
|
||||
@@ -127,6 +127,10 @@ func (b *KafkaBus) Subscribe(ctx context.Context, topics []string, group string,
|
||||
}
|
||||
cc.Set("group.id", group)
|
||||
cc.Set("auto.offset.reset", "earliest")
|
||||
// Offsets are stored once a message is handled and committed in the
|
||||
// background: a synchronous commit per message cost a broker round trip each.
|
||||
cc.Set("enable.auto.commit", true)
|
||||
cc.Set("enable.auto.offset.store", false)
|
||||
|
||||
cons, err := cc.Connect()
|
||||
if err != nil {
|
||||
|
||||
@@ -91,7 +91,10 @@ func (cons *Consumer) Consume(ctx context.Context, handler func(msg *ckf.Message
|
||||
log.Error().Err(err).Msg("kafka message handler error")
|
||||
}
|
||||
|
||||
cons.c.CommitMessage(msg)
|
||||
// Needs enable.auto.offset.store=false; the background commit sends it.
|
||||
if _, err := cons.c.StoreMessage(msg); err != nil {
|
||||
log.Warn().Err(err).Msg("kafka: storing a handled offset failed")
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -12,6 +12,17 @@ type EventWorkerEmailValidation struct {
|
||||
Credentials *SmtpImap `json:"credentials" avro:"credentials"`
|
||||
}
|
||||
|
||||
// EmailValidationRequestChannel is the Redis channel a worker takes credential
|
||||
// checks on, so a check never waits behind the commands queued on its topic.
|
||||
func EmailValidationRequestChannel(workerID uuid.UUID) string {
|
||||
return "email_validation_request:" + workerID.String()
|
||||
}
|
||||
|
||||
// EmailValidationReplyChannel is where the worker answers one check.
|
||||
func EmailValidationReplyChannel(processID uuid.UUID) string {
|
||||
return "email_validation:" + processID.String()
|
||||
}
|
||||
|
||||
// Mail probe reasons: why one leg of a credential check did not pass. They
|
||||
// travel from the worker to the backend inside EmailValidationVerdict and are
|
||||
// what lets the connect form say what the server said instead of "invalid
|
||||
|
||||
@@ -75,6 +75,7 @@ type inboxKitList struct {
|
||||
Username string `json:"username"`
|
||||
Platform string `json:"platform"`
|
||||
Status string `json:"status"`
|
||||
IsAdmin bool `json:"is_admin"`
|
||||
} `json:"mailboxes"`
|
||||
Pages int `json:"pages"`
|
||||
CurrentPage int `json:"current_page"`
|
||||
@@ -128,6 +129,7 @@ func (c *inboxKit) List(ctx context.Context) ([]Mailbox, error) {
|
||||
Provider: normalizeProvider(m.Platform),
|
||||
Status: m.Status,
|
||||
Workspace: ws.Name,
|
||||
Admin: m.IsAdmin,
|
||||
})
|
||||
if len(out) >= MaxMailboxes {
|
||||
return errStop
|
||||
@@ -182,6 +184,72 @@ func (c *inboxKit) show(ctx context.Context, ws, uid, email string) (Credentials
|
||||
return Credentials{Password: res.Password, AppPassword: res.AppPassword}, nil
|
||||
}
|
||||
|
||||
// AuthorizeApp calls POST /v1/api/mailboxes/client-id-request/initiate: a trusted
|
||||
// app under domain-wide delegation on Google, tenant-wide admin consent on Microsoft.
|
||||
func (c *inboxKit) AuthorizeApp(ctx context.Context, a AppAuthorization) (string, error) {
|
||||
ws, _ := unscoped(a.Mailbox.ID)
|
||||
if ws == "" || a.Domain == "" || a.ClientID == "" {
|
||||
return "", invalid(VendorInboxKit, "domain, workspace and client id are required")
|
||||
}
|
||||
body := map[string]any{"domain": a.Domain, "client_id": a.ClientID}
|
||||
switch a.Provider {
|
||||
case ProviderGoogle:
|
||||
body["delegated"], body["scopes"] = true, a.Scopes
|
||||
case ProviderMicrosoft:
|
||||
body["scopes"], body["app_roles"] = a.Scopes, a.AppRoles
|
||||
default:
|
||||
return "", unsupported(VendorInboxKit, "only Google and Microsoft domains can authorize an app")
|
||||
}
|
||||
var res struct {
|
||||
inboxKitEnvelope
|
||||
Data struct {
|
||||
UID string `json:"uid"`
|
||||
} `json:"data"`
|
||||
}
|
||||
if err := c.t.do(ctx, call{method: http.MethodPost, path: "/v1/api/mailboxes/client-id-request/initiate", body: body, header: inboxKitHeader(ws)}, &res); err != nil {
|
||||
return "", err
|
||||
}
|
||||
if err := res.err(); err != nil {
|
||||
return "", err
|
||||
}
|
||||
if res.Data.UID == "" {
|
||||
return "", vendorErr(VendorInboxKit, http.StatusOK, "no request id", nil)
|
||||
}
|
||||
return scoped(ws, res.Data.UID), nil
|
||||
}
|
||||
|
||||
// AuthorizationStatus reads GET /v1/api/mailboxes/client-id-request/status/{id}.
|
||||
func (c *inboxKit) AuthorizationStatus(ctx context.Context, requestID string) (AuthorizationStatus, error) {
|
||||
ws, uid := unscoped(requestID)
|
||||
if ws == "" || uid == "" {
|
||||
return AuthorizationStatus{}, invalid(VendorInboxKit, "request id is required")
|
||||
}
|
||||
var res struct {
|
||||
inboxKitEnvelope
|
||||
Data struct {
|
||||
Status string `json:"status"`
|
||||
ErrorMessage *string `json:"error_message"`
|
||||
} `json:"data"`
|
||||
}
|
||||
if err := c.t.do(ctx, call{method: http.MethodGet, path: "/v1/api/mailboxes/client-id-request/status/" + url.PathEscape(uid), header: inboxKitHeader(ws)}, &res); err != nil {
|
||||
return AuthorizationStatus{}, err
|
||||
}
|
||||
if err := res.err(); err != nil {
|
||||
return AuthorizationStatus{}, err
|
||||
}
|
||||
switch strings.ToLower(res.Data.Status) {
|
||||
case "completed":
|
||||
return AuthorizationStatus{State: AuthorizationCompleted}, nil
|
||||
case "failed", "errored", "cancelled":
|
||||
reason := ""
|
||||
if res.Data.ErrorMessage != nil {
|
||||
reason = *res.Data.ErrorMessage
|
||||
}
|
||||
return AuthorizationStatus{State: AuthorizationFailed, Reason: reason}, nil
|
||||
}
|
||||
return AuthorizationStatus{State: AuthorizationPending}, nil
|
||||
}
|
||||
|
||||
const inboxKitDomainPageSize = 100
|
||||
|
||||
// inboxKitEnvelope carries InboxKit's error flag, which can be set inside an HTTP 200.
|
||||
|
||||
@@ -0,0 +1,85 @@
|
||||
package mailvendor
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"reflect"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestInboxKitAuthorizesAnApp(t *testing.T) {
|
||||
const ws = "6f1c2d3e-0000-4000-8000-000000000001"
|
||||
status := "processing"
|
||||
srv := newRecorder(t, func(w http.ResponseWriter, r *http.Request) {
|
||||
switch r.URL.Path {
|
||||
case "/v1/api/mailboxes/client-id-request/initiate":
|
||||
writeJSON(w, 200, `{"error":false,"message":"Client ID request initiated successfully","data":{"uid":"req-1","status":"queued","domain_name":"acme.io"}}`)
|
||||
case "/v1/api/mailboxes/client-id-request/status/req-1":
|
||||
writeJSON(w, 200, `{"error":false,"message":"ok","data":{"uid":"req-1","status":"`+status+`","error_message":"No admin mailbox found for domain"}}`)
|
||||
default:
|
||||
writeJSON(w, 404, `{}`)
|
||||
}
|
||||
})
|
||||
c := newTestClient(t, VendorInboxKit, map[string]string{FieldAPIKey: testKey}, srv.URL, nil)
|
||||
az, ok := c.(AppAuthorizer)
|
||||
if !ok {
|
||||
t.Fatal("InboxKit client is not an AppAuthorizer")
|
||||
}
|
||||
|
||||
id, err := az.AuthorizeApp(context.Background(), AppAuthorization{
|
||||
Domain: "acme.io", Mailbox: Mailbox{ID: ws + ":m1"}, Provider: ProviderMicrosoft,
|
||||
ClientID: "11111111-2222-3333-4444-555555555555", Scopes: []string{"User.Read"}, AppRoles: []string{"Mail.Send"},
|
||||
})
|
||||
if err != nil || id != ws+":req-1" {
|
||||
t.Fatalf("AuthorizeApp = %q, %v", id, err)
|
||||
}
|
||||
init := srv.requests()[0]
|
||||
assertHeader(t, init, "X-Workspace-Id", ws)
|
||||
var body map[string]any
|
||||
_ = json.Unmarshal([]byte(init.Body), &body)
|
||||
want := map[string]any{
|
||||
"domain": "acme.io", "client_id": "11111111-2222-3333-4444-555555555555",
|
||||
"scopes": []any{"User.Read"}, "app_roles": []any{"Mail.Send"},
|
||||
}
|
||||
if !reflect.DeepEqual(body, want) {
|
||||
t.Fatalf("initiate body = %v\nwant %v", body, want)
|
||||
}
|
||||
|
||||
for _, tc := range []struct{ status, state, reason string }{
|
||||
{"processing", AuthorizationPending, ""},
|
||||
{"completed", AuthorizationCompleted, ""},
|
||||
{"errored", AuthorizationFailed, "No admin mailbox found for domain"},
|
||||
} {
|
||||
status = tc.status
|
||||
got, err := az.AuthorizationStatus(context.Background(), id)
|
||||
if err != nil || got.State != tc.state || got.Reason != tc.reason {
|
||||
t.Fatalf("%s: status = %+v, %v", tc.status, got, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestInboxKitGoogleAuthorizationIsDelegated(t *testing.T) {
|
||||
srv := newRecorder(t, func(w http.ResponseWriter, _ *http.Request) {
|
||||
writeJSON(w, 200, `{"error":false,"data":{"uid":"req-2","status":"queued"}}`)
|
||||
})
|
||||
c := newTestClient(t, VendorInboxKit, map[string]string{FieldAPIKey: testKey}, srv.URL, nil)
|
||||
scopes := []string{"https://www.googleapis.com/auth/gmail.modify"}
|
||||
if _, err := c.(AppAuthorizer).AuthorizeApp(context.Background(), AppAuthorization{
|
||||
Domain: "acme.io", Mailbox: Mailbox{ID: "ws:m1"}, Provider: ProviderGoogle, ClientID: "1234567890", Scopes: scopes, AppRoles: []string{"ignored"},
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var body map[string]any
|
||||
_ = json.Unmarshal([]byte(srv.requests()[0].Body), &body)
|
||||
if body["delegated"] != true || body["app_roles"] != nil || !reflect.DeepEqual(body["scopes"], []any{scopes[0]}) {
|
||||
t.Fatalf("google initiate body = %v", body)
|
||||
}
|
||||
}
|
||||
|
||||
func TestInboxKitAuthorizationNeedsAWorkspace(t *testing.T) {
|
||||
c := newTestClient(t, VendorInboxKit, map[string]string{FieldAPIKey: testKey}, "http://127.0.0.1:1", nil)
|
||||
if _, err := c.(AppAuthorizer).AuthorizeApp(context.Background(), AppAuthorization{Domain: "acme.io", Mailbox: Mailbox{ID: "m1"}, Provider: ProviderMicrosoft, ClientID: "x"}); err == nil {
|
||||
t.Fatal("an id without a workspace was sent")
|
||||
}
|
||||
}
|
||||
@@ -87,6 +87,8 @@ type Mailbox struct {
|
||||
Status string
|
||||
// Workspace is the name of the vendor workspace holding the mailbox, empty where the vendor has none.
|
||||
Workspace string
|
||||
// Admin marks the domain's administrator mailbox, where the vendor says.
|
||||
Admin bool
|
||||
}
|
||||
|
||||
// Endpoint is one server a mailbox connects to.
|
||||
@@ -118,6 +120,38 @@ type Client interface {
|
||||
Credentials(ctx context.Context, m Mailbox) (Credentials, error)
|
||||
}
|
||||
|
||||
// AppAuthorization asks a vendor to authorize an OAuth app on one of its domains,
|
||||
// through the administrator mailbox it holds there.
|
||||
type AppAuthorization struct {
|
||||
Domain string
|
||||
// Mailbox is any listed mailbox on the domain; its id names the vendor workspace.
|
||||
Mailbox Mailbox
|
||||
Provider string
|
||||
ClientID string
|
||||
// Scopes are delegated scopes (Google: full scope URLs); AppRoles are Microsoft application permissions.
|
||||
Scopes []string
|
||||
AppRoles []string
|
||||
}
|
||||
|
||||
// Authorization states.
|
||||
const (
|
||||
AuthorizationPending = "pending"
|
||||
AuthorizationCompleted = "completed"
|
||||
AuthorizationFailed = "failed"
|
||||
)
|
||||
|
||||
// AuthorizationStatus is where an app authorization stands; Reason is the vendor's words on a failure.
|
||||
type AuthorizationStatus struct {
|
||||
State string
|
||||
Reason string
|
||||
}
|
||||
|
||||
// AppAuthorizer is a vendor that can authorize an app on a domain it administers.
|
||||
type AppAuthorizer interface {
|
||||
AuthorizeApp(ctx context.Context, a AppAuthorization) (requestID string, err error)
|
||||
AuthorizationStatus(ctx context.Context, requestID string) (AuthorizationStatus, error)
|
||||
}
|
||||
|
||||
// Option configures New.
|
||||
type Option func(*options)
|
||||
|
||||
|
||||
@@ -211,10 +211,10 @@ type EmailRepository interface {
|
||||
// Used by the warmup reconciler to (re)seed chains.
|
||||
ListWarmupScheduleCandidates(ctx context.Context, limit int) ([]uuid.UUID, error)
|
||||
|
||||
// ListActiveWorkerAccounts returns the ids of every active mailbox. The
|
||||
// worker reconciler uses it to (re)load accounts onto their assigned workers
|
||||
// after onboarding, worker restarts, or reassignment.
|
||||
ListActiveWorkerAccounts(ctx context.Context) ([]uuid.UUID, error)
|
||||
// ListActiveWorkerAccounts returns every active mailbox with the worker it
|
||||
// is assigned to. The worker reconciler uses it to (re)load accounts onto
|
||||
// their assigned workers after onboarding, worker restarts, or reassignment.
|
||||
ListActiveWorkerAccounts(ctx context.Context) ([]MailboxAssignment, error)
|
||||
// ListActiveAccountsByWorker returns the ids of the active mailboxes
|
||||
// assigned to one worker, for reloading them after that worker restarts.
|
||||
ListActiveAccountsByWorker(ctx context.Context, workerID uuid.UUID) ([]uuid.UUID, error)
|
||||
@@ -696,23 +696,29 @@ func (r *emailRepository) ListWarmupScheduleCandidates(ctx context.Context, limi
|
||||
return ids, rows.Err()
|
||||
}
|
||||
|
||||
func (r *emailRepository) ListActiveWorkerAccounts(ctx context.Context) ([]uuid.UUID, error) {
|
||||
const query = `SELECT id FROM email_accounts WHERE status = 'active'`
|
||||
// MailboxAssignment is an active mailbox and the worker it is assigned to, nil when none.
|
||||
type MailboxAssignment struct {
|
||||
ID uuid.UUID
|
||||
WorkerID *uuid.UUID
|
||||
}
|
||||
|
||||
func (r *emailRepository) ListActiveWorkerAccounts(ctx context.Context) ([]MailboxAssignment, error) {
|
||||
const query = `SELECT id, worker_id FROM email_accounts WHERE status = 'active'`
|
||||
rows, err := r.DB.Query(ctx, query)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer rows.Close()
|
||||
|
||||
var ids []uuid.UUID
|
||||
var out []MailboxAssignment
|
||||
for rows.Next() {
|
||||
var id uuid.UUID
|
||||
if err := rows.Scan(&id); err != nil {
|
||||
var a MailboxAssignment
|
||||
if err := rows.Scan(&a.ID, &a.WorkerID); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
ids = append(ids, id)
|
||||
out = append(out, a)
|
||||
}
|
||||
return ids, rows.Err()
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
func (r *emailRepository) ListActiveAccountsByWorker(ctx context.Context, workerID uuid.UUID) ([]uuid.UUID, error) {
|
||||
|
||||
@@ -42,8 +42,17 @@ type ImportWorkRow struct {
|
||||
MailHost string
|
||||
Payload string
|
||||
Attempts int
|
||||
// Code is the row's recorded code; set only by ParkedRows.
|
||||
Code string
|
||||
// AccountID is the mailbox an earlier claim of the row created; set by Claim.
|
||||
AccountID *uuid.UUID
|
||||
// ImportCreatedAt is when the import started; set by Claim.
|
||||
ImportCreatedAt time.Time
|
||||
}
|
||||
|
||||
// ParkedVendorCause is the cause of a row waiting on its vendor to authorize Warmbly; its import is not finished.
|
||||
const ParkedVendorCause = "vendor_authorizing"
|
||||
|
||||
// FinishedImport is an import the last row of which just settled.
|
||||
type FinishedImport struct {
|
||||
ID uuid.UUID
|
||||
@@ -78,6 +87,8 @@ type MailboxImportRepository interface {
|
||||
// FinishRow records a row's outcome. With attempt > 0 it only applies to that claim of the row, so a
|
||||
// replica whose lease lapsed cannot overwrite the outcome another replica recorded.
|
||||
FinishRow(ctx context.Context, id uuid.UUID, line, attempt int, status, code, cause, message string, accountID *uuid.UUID, keepPayload bool) error
|
||||
// SetRowAccount records the mailbox a claimed row just created, before its settings are applied.
|
||||
SetRowAccount(ctx context.Context, id uuid.UUID, line, attempt int, accountID uuid.UUID) error
|
||||
// Touch renews a claimed row's lease when work on it starts; false when another replica has it now.
|
||||
Touch(ctx context.Context, id uuid.UUID, line, attempt int, lease time.Duration) (bool, error)
|
||||
// SettleOrphans closes rows left running in imports that are no longer running.
|
||||
@@ -90,6 +101,12 @@ type MailboxImportRepository interface {
|
||||
// and returns them with the payload they held, so their settings can apply.
|
||||
ResolveSignin(ctx context.Context, orgID uuid.UUID, email string, accountID uuid.UUID) ([]ImportWorkRow, error)
|
||||
PendingSignins(ctx context.Context, limit int) ([]PendingSignin, error)
|
||||
// ParkedRows lists rows waiting on sign-in under one cause that still hold credentials.
|
||||
ParkedRows(ctx context.Context, cause string, limit int) ([]ImportWorkRow, 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
|
||||
// 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)
|
||||
SaveMapping(ctx context.Context, orgID uuid.UUID, signature string, mapping models.MailboxImportMapping) error
|
||||
}
|
||||
@@ -410,6 +427,13 @@ func (r *mailboxImportRepository) Cancel(ctx context.Context, orgID, id uuid.UUI
|
||||
db.CaptureError(err, query, nil, "exec")
|
||||
return err
|
||||
}
|
||||
// Rows parked on a vendor authorization fall back to sign-in, since nothing resumes a cancelled import.
|
||||
query = `UPDATE mailbox_import_rows SET cause = code, message = 'Waiting for someone to sign in as this mailbox.', updated_at = now()
|
||||
WHERE import_id = $1 AND status = 'needs_signin' AND cause = '` + ParkedVendorCause + `'`
|
||||
if _, err := tx.Exec(ctx, query, id); err != nil {
|
||||
db.CaptureError(err, query, nil, "exec")
|
||||
return err
|
||||
}
|
||||
return tx.Commit(ctx)
|
||||
}
|
||||
|
||||
@@ -432,7 +456,7 @@ func (r *mailboxImportRepository) Claim(ctx context.Context, limit int, lease ti
|
||||
SET status = 'running', lease_until = now() + $2::interval, attempts = r.attempts + 1, updated_at = now()
|
||||
FROM due, mailbox_imports i
|
||||
WHERE r.import_id = due.import_id AND r.line = due.line AND i.id = r.import_id
|
||||
RETURNING r.import_id, i.organization_id, i.created_by, i.on_existing, i.settings, r.line, r.email, r.mail_host, r.payload, r.attempts`
|
||||
RETURNING r.import_id, i.organization_id, i.created_by, i.on_existing, i.settings, r.line, r.email, r.mail_host, r.payload, r.attempts, r.email_account_id, i.created_at`
|
||||
rows, err := r.DB.Query(ctx, query, limit, lease.String())
|
||||
if err != nil {
|
||||
db.CaptureError(err, query, nil, "query")
|
||||
@@ -442,7 +466,7 @@ func (r *mailboxImportRepository) Claim(ctx context.Context, limit int, lease ti
|
||||
var out []ImportWorkRow
|
||||
for rows.Next() {
|
||||
var w ImportWorkRow
|
||||
if err := rows.Scan(&w.ImportID, &w.OrgID, &w.CreatedBy, &w.OnExisting, &w.Settings, &w.Line, &w.Email, &w.MailHost, &w.Payload, &w.Attempts); err != nil {
|
||||
if err := rows.Scan(&w.ImportID, &w.OrgID, &w.CreatedBy, &w.OnExisting, &w.Settings, &w.Line, &w.Email, &w.MailHost, &w.Payload, &w.Attempts, &w.AccountID, &w.ImportCreatedAt); err != nil {
|
||||
db.CaptureError(err, query, nil, "scan")
|
||||
return nil, err
|
||||
}
|
||||
@@ -470,6 +494,16 @@ func (r *mailboxImportRepository) FinishRow(ctx context.Context, id uuid.UUID, l
|
||||
return nil
|
||||
}
|
||||
|
||||
func (r *mailboxImportRepository) SetRowAccount(ctx context.Context, id uuid.UUID, line, attempt int, accountID uuid.UUID) error {
|
||||
query := `UPDATE mailbox_import_rows SET email_account_id = $4
|
||||
WHERE import_id = $1 AND line = $2 AND status = 'running' AND attempts = $3`
|
||||
if _, err := r.DB.Exec(ctx, query, id, line, attempt, accountID); err != nil {
|
||||
db.CaptureError(err, query, nil, "exec")
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (r *mailboxImportRepository) Touch(ctx context.Context, id uuid.UUID, line, attempt int, lease time.Duration) (bool, error) {
|
||||
query := `UPDATE mailbox_import_rows SET lease_until = now() + $4::interval
|
||||
WHERE import_id = $1 AND line = $2 AND status = 'running' AND attempts = $3`
|
||||
@@ -503,7 +537,8 @@ func (r *mailboxImportRepository) CompleteFinished(ctx context.Context, credenti
|
||||
WHERE i.status = 'running'
|
||||
AND NOT EXISTS (
|
||||
SELECT 1 FROM mailbox_import_rows r
|
||||
WHERE r.import_id = i.id AND r.status IN ('queued', 'running')
|
||||
WHERE r.import_id = i.id
|
||||
AND (r.status IN ('queued', 'running') OR (r.status = 'needs_signin' AND r.cause = '` + ParkedVendorCause + `'))
|
||||
)
|
||||
RETURNING i.id, i.organization_id, i.created_by`
|
||||
rows, err := r.DB.Query(ctx, query, credentialDays)
|
||||
@@ -602,6 +637,79 @@ func (r *mailboxImportRepository) PendingSignins(ctx context.Context, limit int)
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
func (r *mailboxImportRepository) ParkedRows(ctx context.Context, cause string, limit int) ([]ImportWorkRow, error) {
|
||||
query := `
|
||||
SELECT r.import_id, i.organization_id, i.created_by, r.line, r.email, r.code, r.payload
|
||||
FROM mailbox_import_rows r
|
||||
JOIN mailbox_imports i ON i.id = r.import_id
|
||||
WHERE r.status = 'needs_signin' AND r.cause = $1 AND r.payload <> '' AND i.status = 'running'
|
||||
ORDER BY r.updated_at
|
||||
LIMIT $2`
|
||||
rows, err := r.DB.Query(ctx, query, cause, limit)
|
||||
if err != nil {
|
||||
db.CaptureError(err, query, nil, "query")
|
||||
return nil, err
|
||||
}
|
||||
defer rows.Close()
|
||||
var out []ImportWorkRow
|
||||
for rows.Next() {
|
||||
var w ImportWorkRow
|
||||
if err := rows.Scan(&w.ImportID, &w.OrgID, &w.CreatedBy, &w.Line, &w.Email, &w.Code, &w.Payload); err != nil {
|
||||
db.CaptureError(err, query, nil, "scan")
|
||||
return nil, err
|
||||
}
|
||||
out = append(out, w)
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
func (r *mailboxImportRepository) TouchParked(ctx context.Context, cause string, rows []ImportWorkRow) error {
|
||||
if len(rows) == 0 {
|
||||
return nil
|
||||
}
|
||||
ids, lines := make([]uuid.UUID, len(rows)), make([]int32, len(rows))
|
||||
for i, w := range rows {
|
||||
ids[i], lines[i] = w.ImportID, int32(w.Line)
|
||||
}
|
||||
query := `
|
||||
UPDATE mailbox_import_rows r SET updated_at = now()
|
||||
FROM unnest($2::uuid[], $3::int[]) AS k(import_id, line)
|
||||
WHERE r.import_id = k.import_id AND r.line = k.line AND r.status = 'needs_signin' AND r.cause = $1`
|
||||
if _, err := r.DB.Exec(ctx, query, cause, ids, lines); err != nil {
|
||||
db.CaptureError(err, query, nil, "exec")
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (r *mailboxImportRepository) ResumeParked(ctx context.Context, id uuid.UUID, line int, cause string) error {
|
||||
tx, err := r.DB.Begin(ctx)
|
||||
if err != nil {
|
||||
db.CaptureError(err, "", nil, "begin")
|
||||
return err
|
||||
}
|
||||
defer tx.Rollback(ctx)
|
||||
query := `
|
||||
UPDATE mailbox_import_rows
|
||||
SET status = 'queued', lease_until = NULL, attempts = 0, code = '', cause = '', message = '', updated_at = now()
|
||||
WHERE import_id = $1 AND line = $2 AND status = 'needs_signin' AND cause = $3`
|
||||
tag, err := tx.Exec(ctx, query, id, line, cause)
|
||||
if err != nil {
|
||||
db.CaptureError(err, query, nil, "exec")
|
||||
return err
|
||||
}
|
||||
if tag.RowsAffected() == 0 {
|
||||
return nil
|
||||
}
|
||||
query = `UPDATE mailbox_imports SET status = 'running', finished_at = NULL, credentials_expire_at = NULL, updated_at = now()
|
||||
WHERE id = $1 AND status = 'completed'`
|
||||
if _, err := tx.Exec(ctx, query, id); err != nil {
|
||||
db.CaptureError(err, query, nil, "exec")
|
||||
return err
|
||||
}
|
||||
return tx.Commit(ctx)
|
||||
}
|
||||
|
||||
func (r *mailboxImportRepository) GetMapping(ctx context.Context, orgID uuid.UUID, signature string) (models.MailboxImportMapping, bool, error) {
|
||||
query := `SELECT mapping FROM mailbox_import_mappings WHERE organization_id = $1 AND signature = $2`
|
||||
var raw []byte
|
||||
|
||||
@@ -49,7 +49,15 @@ import buildError from "@/lib/helper/buildError";
|
||||
import timeAgo from "@/lib/helper/timeAgo";
|
||||
import { mailHostLabel, mailHostOAuthProvider } from "@/lib/mailHost";
|
||||
import { cn } from "@/lib/utils";
|
||||
import { ROW_STATUS, importSourceName, isPasswordCause, isSigninCause, plural } from "./importFields";
|
||||
import {
|
||||
AUTHORIZING_STATUS,
|
||||
ROW_STATUS,
|
||||
VENDOR_AUTHORIZING,
|
||||
importSourceName,
|
||||
isPasswordCause,
|
||||
isSigninCause,
|
||||
plural,
|
||||
} from "./importFields";
|
||||
import { Banner, HostMark, Linkified, Pill, SectionLabel, StatCard } from "./parts";
|
||||
import ProviderLogo from "@/components/app/emails/ProviderLogo";
|
||||
import RowFixEditor from "./RowFixEditor";
|
||||
@@ -214,6 +222,9 @@ export default function RunStep({
|
||||
}
|
||||
}
|
||||
|
||||
const authorizing = data.causes.find((x) => x.cause === VENDOR_AUTHORIZING)?.count ?? 0;
|
||||
const signinWaiting = Math.max(0, c.needs_signin - authorizing);
|
||||
const failureCauses = data.causes.filter((x) => x.cause !== VENDOR_AUTHORIZING);
|
||||
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");
|
||||
const reimportHint =
|
||||
@@ -228,9 +239,11 @@ export default function RunStep({
|
||||
? "Import stopped"
|
||||
: c.failed > 0
|
||||
? "Finished with some failures"
|
||||
: c.needs_signin > 0
|
||||
? `Almost done: ${plural(c.needs_signin, "mailbox needs", "mailboxes need")} sign-in`
|
||||
: "All mailboxes imported";
|
||||
: signinWaiting > 0
|
||||
? `Almost done: ${plural(signinWaiting, "mailbox needs", "mailboxes need")} sign-in`
|
||||
: authorizing > 0
|
||||
? `Almost done: ${plural(authorizing, "mailbox is", "mailboxes are")} being authorized`
|
||||
: "All mailboxes imported";
|
||||
|
||||
const page = rows.data;
|
||||
const pageIndex = cursors.length - 1;
|
||||
@@ -243,6 +256,8 @@ export default function RunStep({
|
||||
<Loader2Icon className="w-6 h-6 text-sky-600 animate-spin shrink-0" />
|
||||
) : data.status === "cancelled" ? (
|
||||
<XCircleIcon className="w-6 h-6 text-slate-400 shrink-0" />
|
||||
) : authorizing > 0 && c.failed === 0 && signinWaiting === 0 ? (
|
||||
<Loader2Icon className="w-6 h-6 text-sky-600 animate-spin shrink-0" />
|
||||
) : c.failed > 0 || c.needs_signin > 0 ? (
|
||||
<AlertTriangleIcon className="w-6 h-6 text-amber-600 shrink-0" />
|
||||
) : (
|
||||
@@ -336,12 +351,28 @@ export default function RunStep({
|
||||
</Banner>
|
||||
)}
|
||||
|
||||
{c.needs_signin > 0 && (
|
||||
{authorizing > 0 && (
|
||||
<div className="rounded-md border border-sky-200 bg-sky-50/60 px-3 py-2.5 flex items-start gap-2.5">
|
||||
<Loader2Icon className="w-3.5 h-3.5 text-sky-600 mt-0.5 shrink-0 animate-spin" />
|
||||
<div className="min-w-0 flex-1">
|
||||
<p className="text-[12.5px] font-medium text-sky-900">
|
||||
{vendorLabel(data.vendor) || "Your inbox vendor"} is authorizing Warmbly for{" "}
|
||||
{plural(authorizing, "mailbox", "mailboxes")}
|
||||
</p>
|
||||
<p className="text-[11.5px] text-sky-800/90 leading-relaxed mt-0.5">
|
||||
It approves Warmbly through the admin mailbox it holds on each domain, so nobody has to sign in. This
|
||||
usually takes a few minutes, and the rows connect on their own. You can close this window.
|
||||
</p>
|
||||
</div>
|
||||
</div>
|
||||
)}
|
||||
|
||||
{signinWaiting > 0 && (
|
||||
<div className="rounded-md border border-sky-200 bg-sky-50/60 px-3 py-2.5 flex items-start gap-2.5">
|
||||
<LogInIcon className="w-3.5 h-3.5 text-sky-600 mt-0.5 shrink-0" />
|
||||
<div className="min-w-0 flex-1">
|
||||
<p className="text-[12.5px] font-medium text-sky-900">
|
||||
{plural(c.needs_signin, "mailbox connects", "mailboxes connect")} with Microsoft or Google sign-in
|
||||
{plural(signinWaiting, "mailbox connects", "mailboxes connect")} with Microsoft or Google sign-in
|
||||
</p>
|
||||
<p className="text-[11.5px] text-sky-800/90 leading-relaxed mt-0.5">
|
||||
Their host does not take a password over IMAP. Use Sign in on each row: the window opens for that
|
||||
@@ -384,10 +415,10 @@ export default function RunStep({
|
||||
</div>
|
||||
)}
|
||||
|
||||
{data.causes.length > 0 && (
|
||||
{failureCauses.length > 0 && (
|
||||
<div className="space-y-1.5">
|
||||
<SectionLabel>Why rows failed</SectionLabel>
|
||||
{data.causes.map((cause) => (
|
||||
{failureCauses.map((cause) => (
|
||||
<CauseCard
|
||||
key={cause.cause}
|
||||
cause={cause}
|
||||
@@ -462,7 +493,10 @@ export default function RunStep({
|
||||
<div className="px-3 py-6 text-center text-[11.5px] text-slate-400">No rows here.</div>
|
||||
) : (
|
||||
page.data.map((row) => {
|
||||
const st = ROW_STATUS[row.status] ?? ROW_STATUS.failed;
|
||||
const st =
|
||||
row.status === "needs_signin" && row.cause === VENDOR_AUTHORIZING
|
||||
? AUTHORIZING_STATUS
|
||||
: (ROW_STATUS[row.status] ?? ROW_STATUS.failed);
|
||||
const provider = signInProvider(row);
|
||||
const wantsSignIn = row.status === "needs_signin" || (row.status === "failed" && isSigninCause(row.cause));
|
||||
const canSignIn = wantsSignIn && !!provider;
|
||||
|
||||
@@ -142,6 +142,10 @@ export const ROW_STATUS: Record<ImportRowStatus, { label: string; cls: string }>
|
||||
cancelled: { label: "Cancelled", cls: "text-slate-500 bg-slate-100" },
|
||||
};
|
||||
|
||||
// A needs_signin row parked while the inbox vendor authorizes Warmbly on its domain; it resumes on its own.
|
||||
export const VENDOR_AUTHORIZING = "vendor_authorizing";
|
||||
export const AUTHORIZING_STATUS = { label: "Authorizing", cls: "text-sky-700 bg-sky-50" };
|
||||
|
||||
// Cause keys from the backend's mailcause catalogue whose fix is a new password.
|
||||
const PASSWORD_CAUSE = /(_app_password_required|_bad_credentials)$|^(auth_refused|missing_password)$/;
|
||||
// Cause keys whose fix is a provider sign-in rather than any password; retrying them changes nothing.
|
||||
|
||||
Reference in New Issue
Block a user