feat: carry a unibox read or unread change out to the mailbox itself through a new MESSAGE_SEEN worker command, so a conversation read in Warmbly stops showing bold in Gmail, Outlook and IMAP, relaying the state the row holds rather than the one the request asked for, only for messages that actually changed, dispatched detached from the request and never retried (#515)

This commit is contained in:
Matthew Meszaros
2026-09-14 21:20:42 -07:00
committed by GitHub
parent 20a56dc41b
commit 1ca3bd19b3
19 changed files with 792 additions and 22 deletions
+3
View File
@@ -1356,6 +1356,9 @@ func main() {
// tasksClient isn't initialised until the Cloud Tasks config
// block runs.
uniboxService = unibox.NewService(cache, s3, uniboxRepository, taskRepository, tasksClient)
// Read and unread in the unibox are carried out to the mailbox itself,
// so a thread read here is read in Gmail too.
uniboxService.WireProviderRelay(eventsPublisher)
// Org AI skills (playbooks): CRUD for settings + prompt injection + the
// load_skill tool source.
@@ -258,6 +258,8 @@ Returns the resulting label set in a `data` array.
Marks messages as read or unread, org-wide: either an explicit batch of up to 500 ids, or a whole canonical folder at once. Send one of `email_ids` or `folder`; sending both returns `400`. Auth: **Scope** `WRITE_UNIBOX` · **Org permission** `access_unibox`.
The change is also carried out to the mailbox itself, so a message marked read here stops showing as unread in Gmail, Outlook or an IMAP mailbox. Opening a message with `GET /unibox/:id` marks it read the same way, relay included. Only messages whose state actually changed are relayed, and the relay is best-effort: it happens after the response, on the worker holding the mailbox, and a mailbox that is unplaced at that moment keeps its own read state until something changes it again. Unlike [filing a conversation](/guides/unibox/#filing-a-conversation), which stays inside Warmbly.
### Request body
| Field | Type | Required | Description |
+3 -1
View File
@@ -48,7 +48,9 @@ The backend drives workers with a `{type, body}` envelope (`internal/models/even
}
```
Types: `SEND_EMAIL`, `ADD_EMAIL`, `REMOVE_EMAIL`, `EMAIL_VALIDATION`, `WARMUP_ACTION`. Each has a typed body struct in `internal/models/` (for example `models.SendEmail` for `SEND_EMAIL`).
Types: `SEND_EMAIL`, `ADD_EMAIL`, `REMOVE_EMAIL`, `EMAIL_VALIDATION`, `WARMUP_ACTION`, `MESSAGE_SEEN`. Each has a typed body struct in `internal/models/` (for example `models.SendEmail` for `SEND_EMAIL`).
`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.
## Worker results
+6 -2
View File
@@ -106,7 +106,11 @@ Categories label conversations. Tags label the mailboxes themselves (grouping ac
## Read state
A conversation is unread when any message inside is unseen, marked by a blue dot in the margin, bolder text, and a blue timestamp. Opening it marks its messages seen and updates the counts. Read state syncs both ways with the mailbox provider.
A conversation is unread when any message inside is unseen, marked by a blue dot in the margin, bolder text, and a blue timestamp. Opening it marks its messages seen and updates the counts.
Read state travels both ways. A message arriving already read is stored that way, and reading a conversation here marks it read in the mailbox itself, so Gmail, Outlook and any IMAP mailbox stop showing it bold too. **Mark as unread** does the same in reverse, and **Mark all as read** on a folder carries the whole sweep across. Only a real change travels: opening a conversation you have already read costs nothing at the provider.
The relay runs on the worker that holds the mailbox, just after Warmbly's own copy is updated, so the dashboard never waits on it. If the mailbox is between workers at that moment, or the provider refuses, Warmbly keeps your read state and the mailbox keeps its own; nothing retries, because the next sync reports whatever the provider actually thinks.
## Filing a conversation
@@ -121,7 +125,7 @@ The thread header carries three filing actions, on the row above the message on
Archive and Delete both offer **Undo** on the confirmation toast. Open the Trash or Archive folder and the same header offers **Move to inbox**, so nothing filed by accident is stuck.
<Callout type="info" title="Filing here does not move the message at the provider">
Archive and Delete are Warmbly's own filing. The message keeps its place in Gmail, Outlook, or whatever mail client the mailbox belongs to, and deleting a conversation here never deletes mail there. The next sync will not undo your filing either: Warmbly records where the provider has each message separately from where you filed it, and follows the provider only when the provider itself moves the message. So junking a message in Gmail still reaches Warmbly, and an ordinary sync pass does not.
Unlike read state, Archive and Delete are Warmbly's own filing. The message keeps its place in Gmail, Outlook, or whatever mail client the mailbox belongs to, and deleting a conversation here never deletes mail there. The next sync will not undo your filing either: Warmbly records where the provider has each message separately from where you filed it, and follows the provider only when the provider itself moves the message. So junking a message in Gmail still reaches Warmbly, and an ordinary sync pass does not.
</Callout>
## Replying
+11
View File
@@ -32,6 +32,17 @@ func (s *uniboxService) GetByID(
}
ownerID = owner
accountID = msg.EmailID
// Opening a conversation is what marks it read, org-scoped so any
// member clears the shared unread state. Relayed like any other read
// state change, so the mailbox agrees: this is the path a developer
// hitting the API reaches, where nothing calls PATCH /unibox/seen.
if !msg.Seen {
if changed, err := s.uniboxRepository.MarkSeenBulk(ctx, orgID, []uuid.UUID{id}, true); err == nil {
msg.Seen = true
s.relaySeen(ctx, orgID, changed)
}
}
resp.ID = msg.ID
resp.GmailID = msg.GmailID
resp.UID = msg.UID
+90 -2
View File
@@ -2,8 +2,10 @@ package unibox
import (
"context"
"time"
"github.com/google/uuid"
"github.com/rs/zerolog/log"
"github.com/warmbly/warmbly/internal/errx"
"github.com/warmbly/warmbly/internal/models"
"github.com/warmbly/warmbly/internal/observability/errs"
@@ -32,21 +34,107 @@ func (s *uniboxService) MarkSeenBulk(ctx context.Context, orgID uuid.UUID, data
if !models.ValidFolder(data.Folder) {
return nil, errx.ErrUniboxFolder
}
if err := s.uniboxRepository.MarkSeenByFolder(ctx, orgID, data.Folder, data.Seen); err != nil {
changed, err := s.uniboxRepository.MarkSeenByFolder(ctx, orgID, data.Folder, data.Seen)
if err != nil {
errs.CaptureException(err)
return nil, errx.InternalError()
}
s.relaySeen(ctx, orgID, changed)
return data, nil
}
if err := s.uniboxRepository.MarkSeenBulk(ctx, orgID, data.EmailIDs, data.Seen); err != nil {
changed, err := s.uniboxRepository.MarkSeenBulk(ctx, orgID, data.EmailIDs, data.Seen)
if err != nil {
errs.CaptureException(err)
return nil, errx.InternalError()
}
s.relaySeen(ctx, orgID, changed)
return data, nil
}
// relaySeen carries a read/unread change out to the mailboxes themselves, so
// a thread read in Warmbly is read in Gmail too.
//
// Detached and best-effort. The store is the customer's view and has already
// been written, so this must not hold the response open behind a slow broker,
// and a worker that cannot be reached must not fail the request. Nothing
// retries: the provider's own state is what the next sync brings back anyway.
//
// What is relayed is the state the ROW now holds, read back inside the
// lookup, not the state the request asked for. Two people toggling the same
// conversation in opposite directions at the same moment can still race, but
// the loser then relays the winner's answer rather than its own.
func (s *uniboxService) relaySeen(ctx context.Context, orgID uuid.UUID, changed []uuid.UUID) {
if s.publisher == nil || len(changed) == 0 {
return
}
go func() {
// Detached from the request, bounded so a wedged broker cannot leak a
// goroutine per press.
bg, cancel := context.WithTimeout(context.WithoutCancel(ctx), seenRelayTimeout)
defer cancel()
s.publishSeenRelay(bg, orgID, changed)
}()
}
// seenRelayTimeout bounds one relay. Generous, because "mark all as read" on
// a busy folder is many events, and it exists to end a wedged publish rather
// than to pace a healthy one.
const seenRelayTimeout = 2 * time.Minute
func (s *uniboxService) publishSeenRelay(ctx context.Context, orgID uuid.UUID, changed []uuid.UUID) {
targets, err := s.uniboxRepository.SeenRelayTargets(ctx, orgID, changed)
if err != nil {
errs.CaptureException(err)
return
}
// One event per mailbox and read state: "mark all as read" on a busy
// folder is one press over thousands of messages, a mailbox is the unit a
// worker holds, and a batch carries one state for all of it.
type relayKey struct {
emailID uuid.UUID
seen bool
}
batches := make(map[relayKey]*models.MessageSeenAction)
workers := make(map[uuid.UUID]uuid.UUID, len(targets))
order := make([]relayKey, 0, len(targets))
for _, t := range targets {
key := relayKey{emailID: t.EmailID, seen: t.Seen}
act, ok := batches[key]
if !ok {
act = &models.MessageSeenAction{EmailID: t.EmailID, Seen: t.Seen}
batches[key] = act
workers[t.EmailID] = t.WorkerID
order = append(order, key)
}
act.Messages = append(act.Messages, t.Ref)
}
for _, key := range order {
act := batches[key]
for start := 0; start < len(act.Messages); start += models.SeenRelayChunk {
end := start + models.SeenRelayChunk
if end > len(act.Messages) {
end = len(act.Messages)
}
batch := &models.MessageSeenAction{
EmailID: key.emailID,
Seen: key.seen,
Messages: act.Messages[start:end],
}
if err := s.publisher.PublishMessageSeen(ctx, workers[key.emailID], batch); err != nil {
log.Warn().Err(err).
Str("email_account_id", key.emailID.String()).
Int("messages", len(batch.Messages)).
Msg("could not relay the unibox read state to the mailbox provider")
}
}
}
}
// MoveFolderBulk backs Archive, Delete and Move to inbox in the thread header.
// Store-side only: the provider copy stays where it is, and provider_folder is
// left alone so the sync can still tell a real provider move from a flag scan.
+166
View File
@@ -0,0 +1,166 @@
package unibox
import (
"context"
"os"
"testing"
"github.com/google/uuid"
"github.com/jackc/pgx/v5/pgxpool"
"github.com/warmbly/warmbly/internal/infrastructure/db"
"github.com/warmbly/warmbly/internal/repository"
)
// Live cover for what the relay is built on: that marking read reports only
// the messages that actually changed, and that those resolve to the provider
// handles and the worker holding the mailbox. Skipped unless WARMBLY_TEST_DB
// is set:
//
// WARMBLY_TEST_DB=postgres://warmbly:warmbly@localhost:15432/warmbly_dev?sslmode=disable \
// go test ./internal/app/unibox/ -run Live -v
type seenLiveFixture struct {
pool *pgxpool.Pool
repo repository.UniboxRepository
user uuid.UUID
org uuid.UUID
worker uuid.UUID
box uuid.UUID
unread uuid.UUID
read uuid.UUID
}
func newSeenLiveFixture(t *testing.T) *seenLiveFixture {
t.Helper()
dsn := os.Getenv("WARMBLY_TEST_DB")
if dsn == "" {
t.Skip("WARMBLY_TEST_DB not set")
}
handle, err := db.New(context.Background(), dsn)
if err != nil {
t.Fatalf("connect: %v", err)
}
t.Cleanup(func() { handle.Pool.Close() })
ctx := context.Background()
f := &seenLiveFixture{
pool: handle.Pool, user: uuid.New(), org: uuid.New(),
worker: uuid.New(), box: uuid.New(), unread: uuid.New(), read: uuid.New(),
}
exec := func(sql string, args ...any) {
t.Helper()
if _, err := f.pool.Exec(ctx, sql, args...); err != nil {
t.Fatalf("fixture: %v", err)
}
}
exec(`INSERT INTO users (id, email, first_name, last_name) VALUES ($1, $2, 'Seen', 'Test')`,
f.user, "seen-"+f.user.String()[:8]+"@test.local")
exec(`INSERT INTO organizations (id, name, slug, owner_user_id) VALUES ($1, 'Seen Test', $2, $3)`,
f.org, "seen-"+f.org.String()[:8], f.user)
exec(`INSERT INTO fleet_nodes (id, role, name, address, active, last_seen_at)
VALUES ($1, 'worker', 'seen-test', '127.0.0.1', true, now())`, f.worker)
exec(`INSERT INTO workers (id, account_count, load_score) VALUES ($1, 0, 0)`, f.worker)
exec(`INSERT INTO email_accounts (id, user_id, organization_id, worker_id, email, name,
signature_plain, signature_html, provider, status, campaign_limit, min_wait_time)
VALUES ($1, $2, $3, $4, $5, 'Seen', '', '', 'gmail', 'active', 50, 600)`,
f.box, f.user, f.org, f.worker, "seen-"+f.box.String()[:8]+"@test.local")
exec(`INSERT INTO unibox_emails (id, user_id, email_id, gmail_id, uid, folder_path, message_id, seen, folder)
VALUES ($1, $2, $3, 'gmail-unread', 11, 'INBOX', '<unread@test>', false, 'inbox')`,
f.unread, f.user, f.box)
exec(`INSERT INTO unibox_emails (id, user_id, email_id, gmail_id, uid, folder_path, message_id, seen, folder)
VALUES ($1, $2, $3, 'gmail-read', 12, 'INBOX', '<read@test>', true, 'inbox')`,
f.read, f.user, f.box)
t.Cleanup(func() {
c := context.Background()
for _, step := range []struct {
sql string
arg any
}{
{`DELETE FROM unibox_emails WHERE email_id = $1`, f.box},
{`DELETE FROM email_accounts WHERE id = $1`, f.box},
{`DELETE FROM fleet_nodes WHERE id = $1`, f.worker},
{`DELETE FROM organizations WHERE id = $1`, f.org},
{`DELETE FROM users WHERE id = $1`, f.user},
} {
if _, err := f.pool.Exec(c, step.sql, step.arg); err != nil {
t.Errorf("cleanup %q: %v", step.sql, err)
}
}
})
f.repo = repository.NewUniboxRepository(handle)
return f
}
func TestLiveMarkSeenReportsOnlyRealChanges(t *testing.T) {
f := newSeenLiveFixture(t)
ctx := context.Background()
// Both are asked for; only the unread one is a change, and re-reading a
// thread that is already read must cost nothing at the provider.
changed, err := f.repo.MarkSeenBulk(ctx, f.org, []uuid.UUID{f.unread, f.read}, true)
if err != nil {
t.Fatalf("mark seen: %v", err)
}
if len(changed) != 1 || changed[0] != f.unread {
t.Fatalf("expected only the unread message to be reported, got %v", changed)
}
// A second pass changes nothing at all.
changed, err = f.repo.MarkSeenBulk(ctx, f.org, []uuid.UUID{f.unread, f.read}, true)
if err != nil {
t.Fatalf("mark seen again: %v", err)
}
if len(changed) != 0 {
t.Errorf("a repeated mark-read reported %d changes", len(changed))
}
// Marking unread is a change again, in the other direction.
changed, err = f.repo.MarkSeenBulk(ctx, f.org, []uuid.UUID{f.unread}, false)
if err != nil || len(changed) != 1 {
t.Fatalf("mark unread: %v / %v", err, changed)
}
}
func TestLiveSeenRelayTargets(t *testing.T) {
f := newSeenLiveFixture(t)
ctx := context.Background()
targets, err := f.repo.SeenRelayTargets(ctx, f.org, []uuid.UUID{f.unread, f.read})
if err != nil {
t.Fatalf("relay targets: %v", err)
}
if len(targets) != 2 {
t.Fatalf("expected both messages, got %d", len(targets))
}
for _, target := range targets {
if target.EmailID != f.box || target.WorkerID != f.worker {
t.Errorf("wrong routing: mailbox=%v worker=%v", target.EmailID, target.WorkerID)
}
// All three providers' handles travel; the worker takes what it uses.
if target.Ref.ProviderID == "" || target.Ref.UID == 0 || target.Ref.Folder != "INBOX" || target.Ref.RFCMessageID == "" {
t.Errorf("incomplete ref: %+v", target.Ref)
}
// The state comes from the row, so the two fixtures disagree.
if target.Ref.ProviderID == "gmail-read" && !target.Seen {
t.Error("a read message reported unread")
}
if target.Ref.ProviderID == "gmail-unread" && target.Seen {
t.Error("an unread message reported read")
}
}
// A mailbox with no worker has nothing to relay through and must not
// produce a target keyed on a nil worker.
if _, err := f.pool.Exec(ctx, `UPDATE email_accounts SET worker_id = NULL WHERE id = $1`, f.box); err != nil {
t.Fatalf("unassign worker: %v", err)
}
targets, err = f.repo.SeenRelayTargets(ctx, f.org, []uuid.UUID{f.unread, f.read})
if err != nil {
t.Fatalf("relay targets without a worker: %v", err)
}
if len(targets) != 0 {
t.Errorf("expected no targets for an unplaced mailbox, got %d", len(targets))
}
}
+147
View File
@@ -0,0 +1,147 @@
package unibox
import (
"context"
"testing"
"github.com/google/uuid"
"github.com/warmbly/warmbly/internal/events"
"github.com/warmbly/warmbly/internal/models"
"github.com/warmbly/warmbly/internal/repository"
)
// Only the two methods the relay uses are implemented; the embedded interface
// satisfies the rest and panics if the relay ever reaches for one.
type fakeSeenRepo struct {
repository.UniboxRepository
targets []models.SeenRelayTarget
asked [][]uuid.UUID
}
func (f *fakeSeenRepo) SeenRelayTargets(_ context.Context, _ uuid.UUID, ids []uuid.UUID) ([]models.SeenRelayTarget, error) {
f.asked = append(f.asked, ids)
return f.targets, nil
}
type fakeSeenPublisher struct {
events.Publisher
sent []*models.MessageSeenAction
toWork []uuid.UUID
}
func (f *fakeSeenPublisher) PublishMessageSeen(_ context.Context, workerID uuid.UUID, action *models.MessageSeenAction) error {
f.sent = append(f.sent, action)
f.toWork = append(f.toWork, workerID)
return nil
}
func TestRelaySeenGroupsByMailbox(t *testing.T) {
boxA, boxB := uuid.New(), uuid.New()
workerA, workerB := uuid.New(), uuid.New()
repo := &fakeSeenRepo{targets: []models.SeenRelayTarget{
{EmailID: boxA, WorkerID: workerA, Seen: true, Ref: models.MessageSeenRef{ProviderID: "a1"}},
{EmailID: boxA, WorkerID: workerA, Seen: true, Ref: models.MessageSeenRef{ProviderID: "a2"}},
{EmailID: boxB, WorkerID: workerB, Seen: true, Ref: models.MessageSeenRef{UID: 7, Folder: "INBOX"}},
}}
pub := &fakeSeenPublisher{}
s := &uniboxService{uniboxRepository: repo, publisher: pub}
s.publishSeenRelay(context.Background(), uuid.New(), []uuid.UUID{uuid.New(), uuid.New(), uuid.New()})
if len(pub.sent) != 2 {
t.Fatalf("expected one event per mailbox, got %d", len(pub.sent))
}
// A mailbox is what a worker holds, so each event must go to its own.
for i, act := range pub.sent {
switch act.EmailID {
case boxA:
if pub.toWork[i] != workerA || len(act.Messages) != 2 {
t.Errorf("mailbox A: worker=%v messages=%d", pub.toWork[i], len(act.Messages))
}
case boxB:
if pub.toWork[i] != workerB || len(act.Messages) != 1 {
t.Errorf("mailbox B: worker=%v messages=%d", pub.toWork[i], len(act.Messages))
}
default:
t.Errorf("unexpected mailbox %v", act.EmailID)
}
if !act.Seen {
t.Error("the requested state did not travel")
}
}
}
// "Mark all as read" on a busy folder is one press over thousands of
// messages; the bus must not see one enormous event.
func TestRelaySeenChunks(t *testing.T) {
box, worker := uuid.New(), uuid.New()
targets := make([]models.SeenRelayTarget, 0, models.SeenRelayChunk+100)
for i := 0; i < models.SeenRelayChunk+100; i++ {
targets = append(targets, models.SeenRelayTarget{EmailID: box, WorkerID: worker, Ref: models.MessageSeenRef{ProviderID: "m"}})
}
pub := &fakeSeenPublisher{}
s := &uniboxService{uniboxRepository: &fakeSeenRepo{targets: targets}, publisher: pub}
s.publishSeenRelay(context.Background(), uuid.New(), []uuid.UUID{uuid.New()})
if len(pub.sent) != 2 {
t.Fatalf("expected 2 chunks, got %d", len(pub.sent))
}
if len(pub.sent[0].Messages) != models.SeenRelayChunk || len(pub.sent[1].Messages) != 100 {
t.Errorf("chunk sizes %d and %d", len(pub.sent[0].Messages), len(pub.sent[1].Messages))
}
if pub.sent[0].Seen || pub.sent[1].Seen {
t.Error("marking unread must travel as unread")
}
}
// Nothing changed means nothing to tell the provider, and the repository is
// not even asked.
func TestRelaySeenSkipsWhenNothingChanged(t *testing.T) {
repo := &fakeSeenRepo{}
pub := &fakeSeenPublisher{}
s := &uniboxService{uniboxRepository: repo, publisher: pub}
s.relaySeen(context.Background(), uuid.New(), nil)
if len(pub.sent) != 0 || len(repo.asked) != 0 {
t.Errorf("relayed %d events after %d lookups for no change", len(pub.sent), len(repo.asked))
}
}
// An install with no worker bus wired still has a working unibox.
func TestRelaySeenWithoutAPublisher(t *testing.T) {
s := &uniboxService{uniboxRepository: &fakeSeenRepo{}}
s.relaySeen(context.Background(), uuid.New(), []uuid.UUID{uuid.New()})
}
// A relay batch carries one state, so rows that disagree (a concurrent toggle
// landed between the write and the lookup) have to split rather than be sent
// under whichever state was asked for.
func TestRelaySeenSplitsOnStoredState(t *testing.T) {
box, worker := uuid.New(), uuid.New()
repo := &fakeSeenRepo{targets: []models.SeenRelayTarget{
{EmailID: box, WorkerID: worker, Seen: true, Ref: models.MessageSeenRef{ProviderID: "read"}},
{EmailID: box, WorkerID: worker, Seen: false, Ref: models.MessageSeenRef{ProviderID: "unread"}},
}}
pub := &fakeSeenPublisher{}
s := &uniboxService{uniboxRepository: repo, publisher: pub}
s.publishSeenRelay(context.Background(), uuid.New(), []uuid.UUID{uuid.New(), uuid.New()})
if len(pub.sent) != 2 {
t.Fatalf("expected one event per state, got %d", len(pub.sent))
}
for _, act := range pub.sent {
if len(act.Messages) != 1 {
t.Fatalf("states were merged into one batch: %+v", act)
}
want := "unread"
if act.Seen {
want = "read"
}
if act.Messages[0].ProviderID != want {
t.Errorf("message %q relayed with seen=%v", act.Messages[0].ProviderID, act.Seen)
}
}
}
+17
View File
@@ -6,6 +6,7 @@ import (
"github.com/google/uuid"
"github.com/warmbly/warmbly/internal/errx"
"github.com/warmbly/warmbly/internal/events"
"github.com/warmbly/warmbly/internal/infrastructure/cache"
"github.com/warmbly/warmbly/internal/infrastructure/storage"
"github.com/warmbly/warmbly/internal/models"
@@ -79,6 +80,10 @@ type UniboxService interface {
// before bodies were indexed. Runs until the archive is caught up, then
// returns; blocking, so callers run it in a goroutine.
StartBodyTextBackfill(ctx context.Context)
// WireProviderRelay attaches the worker bus, after which a read/unread
// change made here is carried out to the mailbox provider too.
WireProviderRelay(p events.Publisher)
}
type uniboxService struct {
@@ -87,6 +92,18 @@ type uniboxService struct {
tasksClient tasksched.Scheduler
cache *cache.Cache
blob storage.Store
// publisher relays read/unread changes out to the mailbox providers.
// Optional: without it the unibox still works and only Warmbly's own copy
// of the read state changes.
publisher events.Publisher
}
// WireProviderRelay attaches the bus the unibox relays read state through.
// Wired after construction, like the webhook dispatcher on the mailbox
// service, because a deployment without a worker bus is still a working
// unibox.
func (s *uniboxService) WireProviderRelay(p events.Publisher) {
s.publisher = p
}
func NewService(
+108
View File
@@ -0,0 +1,108 @@
package worker
import (
"context"
"github.com/rs/zerolog/log"
"github.com/warmbly/warmbly/internal/app/worker/wmail"
"github.com/warmbly/warmbly/internal/models"
)
// HandleMessageSeen applies a read/unread change made in the unibox to the
// mailbox itself, so the provider's copy agrees with what the customer sees.
//
// Best-effort by design: the store is already updated when this arrives, and a
// message the provider has since moved or deleted must not turn into a failed
// job. Every provider path logs and moves on, and the next sync brings back
// whatever the provider actually thinks.
func (w *WorkerService) HandleMessageSeen(ctx context.Context, action models.MessageSeenAction) error {
if len(action.Messages) == 0 {
return nil
}
w.mailManager.RLock()
mail, exists := w.mailManager.Emails[action.EmailID]
w.mailManager.RUnlock()
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
// provider's copy misses out.
log.Warn().Str("email_id", action.EmailID.String()).Msg("Mailbox not loaded; read state not relayed to the provider")
return nil
}
log.Info().
Str("email_id", action.EmailID.String()).
Bool("seen", action.Seen).
Int("messages", len(action.Messages)).
Msg("Relaying unibox read state to the provider")
switch {
case mail.GoogleData != nil && mail.GoogleData.Client != nil:
w.relaySeenGoogle(ctx, mail, action)
case mail.GraphData != nil && mail.GraphData.Client != nil:
w.relaySeenGraph(ctx, mail, action)
case mail.SmtpImapData != nil && mail.SmtpImapData.ImapClient != nil:
w.relaySeenImap(ctx, mail, action)
default:
log.Warn().Str("email_id", action.EmailID.String()).Msg("No mail client available to relay read state")
}
return nil
}
// relaySeenGoogle sends the whole batch as one batchModify.
func (w *WorkerService) relaySeenGoogle(ctx context.Context, mail *wmail.WMail, action models.MessageSeenAction) {
ids := make([]string, 0, len(action.Messages))
for _, m := range action.Messages {
if m.ProviderID != "" {
ids = append(ids, m.ProviderID)
}
}
if len(ids) == 0 {
return
}
if err := mail.GoogleData.Client.SetSeen(ctx, ids, action.Seen); err != nil {
log.Warn().Err(err).Str("email_id", action.EmailID.String()).Int("messages", len(ids)).
Msg("Failed to relay read state to Gmail")
}
}
// relaySeenGraph patches one message at a time, re-resolving ids that moved.
func (w *WorkerService) relaySeenGraph(ctx context.Context, mail *wmail.WMail, action models.MessageSeenAction) {
client := mail.GraphData.Client
for _, m := range action.Messages {
id := m.ProviderID
// A Graph message id changes when the message is moved (copy+delete),
// so prefer the immutable Message-ID when one travelled.
if m.RFCMessageID != "" {
if resolved, err := client.ResolveMessageID(ctx, m.RFCMessageID); err == nil && resolved != "" {
id = resolved
}
}
if id == "" {
continue
}
if err := client.SetSeen(ctx, id, action.Seen); err != nil {
log.Warn().Err(err).Str("email_id", action.EmailID.String()).Str("graph_id", id).
Msg("Failed to relay read state to Outlook")
}
}
}
// relaySeenImap groups by folder, because a UID only means anything inside
// one and each folder costs a SELECT.
func (w *WorkerService) relaySeenImap(ctx context.Context, mail *wmail.WMail, action models.MessageSeenAction) {
byFolder := make(map[string][]uint32)
for _, m := range action.Messages {
if m.UID == 0 || m.Folder == "" {
continue
}
byFolder[m.Folder] = append(byFolder[m.Folder], m.UID)
}
for folder, uids := range byFolder {
if err := mail.SmtpImapData.ImapClient.SetSeen(ctx, folder, uids, action.Seen); err != nil {
log.Warn().Err(err).Str("email_id", action.EmailID.String()).Str("folder", folder).Int("uids", len(uids)).
Msg("Failed to relay read state over IMAP")
}
}
}
+1
View File
@@ -26,6 +26,7 @@ func (w *WorkerService) InitEvents() {
Register(w, models.WorkerEventTypeRemoveEmail, w.HandleRemoveEmail)
Register(w, models.WorkerEventTypeEmailValidation, w.HandleEmailValidation)
Register(w, models.WorkerEventTypeWarmupAction, w.HandleWarmupAction)
Register(w, models.WorkerEventTypeMessageSeen, w.HandleMessageSeen)
}
func Register[T any](w *WorkerService, eventType models.WorkerEventType, handler EventHandler[T]) {
+3
View File
@@ -38,6 +38,9 @@ type ImapConn interface {
// Warmup actions.
MarkAsRead(ctx context.Context, mailboxName string, uid uint32) error
// SetSeen is the unibox's read/unread relay: many UIDs in one folder, in
// one STORE, in either direction.
SetSeen(ctx context.Context, mailboxName string, uids []uint32, seen bool) error
MarkImportant(ctx context.Context, mailboxName string, uid uint32) error
MoveToFolder(ctx context.Context, sourceMailbox, dstFolder string, uid uint32) error
RemoveFromSpam(ctx context.Context, sourceMailbox, inboxName string, uid uint32) error
+26
View File
@@ -122,3 +122,29 @@ func (c *Client) getOrCreateLabel(ctx context.Context, labelName string) (string
return newLabel.Id, nil
}
// SetSeen flips the read state of many messages in one call.
//
// batchModify rather than a modify per message: the unibox's "mark all as
// read" is one press over a folder, and a per-message call there is hundreds
// of round trips against a per-user rate limit. Gmail accepts up to 1000 ids
// per request; the caller chunks.
func (c *Client) SetSeen(ctx context.Context, messageIDs []string, seen bool) error {
if c.srv == nil {
return fmt.Errorf("gmail service not initialized")
}
if len(messageIDs) == 0 {
return nil
}
req := &gmail.BatchModifyMessagesRequest{Ids: messageIDs}
if seen {
req.RemoveLabelIds = []string{Unread}
} else {
req.AddLabelIds = []string{Unread}
}
if err := c.srv.Users.Messages.BatchModify("me", req).Context(ctx).Do(); err != nil {
return fmt.Errorf("failed to set read state on %d messages: %w", len(messageIDs), err)
}
return nil
}
+7
View File
@@ -182,3 +182,10 @@ func (c *Client) cacheFolder(name, id string) {
c.folderIDs[name] = id
c.mu.Unlock()
}
// SetSeen flips the read state of one message. Graph has no batch equivalent
// of Gmail's batchModify that is worth the complexity here, so the caller
// loops.
func (c *Client) SetSeen(ctx context.Context, messageID string, seen bool) error {
return c.doJSON(ctx, "PATCH", c.messageURL(messageID), map[string]any{"isRead": seen}, nil)
}
@@ -217,3 +217,47 @@ func IsInboxMailbox(name string, attrs []string) bool {
}
return false
}
// SetSeen adds or removes \Seen across many UIDs in one STORE.
//
// One command for the whole set, because the alternative is a SELECT and a
// STORE per message and the unibox files them in bulk. UIDs the server no
// longer has are silently ignored by STORE, which is what should happen: the
// message moved or was deleted at the provider and the next sync will say so.
func (c *Client) SetSeen(ctx context.Context, mailboxName string, uids []uint32, seen bool) error {
if len(uids) == 0 {
return nil
}
c.mu.Lock()
defer c.mu.Unlock()
if merr := c.ensureConnected(); merr != nil {
return merr
}
c.lifecycle.RLock()
defer c.lifecycle.RUnlock()
defer c.begin()()
if _, err := c.selectMailbox(mailboxName, nil); err != nil {
return fmt.Errorf("select %q: %w", mailboxName, err)
}
set := imap.UIDSet{}
for _, uid := range uids {
set.AddNum(imap.UID(uid))
}
op := imap.StoreFlagsDel
if seen {
op = imap.StoreFlagsAdd
}
cmd := c.client.Store(set, &imap.StoreFlags{
Op: op,
Silent: true,
Flags: []imap.Flag{imap.FlagSeen},
}, nil)
if err := cmd.Close(); err != nil {
return fmt.Errorf("store \\Seen (%d uids): %w", len(uids), err)
}
return nil
}
+16
View File
@@ -35,6 +35,9 @@ type Publisher interface {
// Warmup action events
PublishWarmupAction(ctx context.Context, workerID uuid.UUID, action *models.WarmupEmailAction) error
// PublishMessageSeen relays a read/unread change made in the unibox out to
// the mailbox provider.
PublishMessageSeen(ctx context.Context, workerID uuid.UUID, action *models.MessageSeenAction) error
// Worker change notifications
PublishAddEmail(ctx context.Context, workerID uuid.UUID, email *models.AddWorkerEmail) error
@@ -310,6 +313,19 @@ func (p *publisher) PublishWarmupAction(ctx context.Context, workerID uuid.UUID,
return p.publish(workerTopic, action.EmailID.String(), workerEvent)
}
// PublishMessageSeen relays a unibox read/unread change to the worker holding
// the mailbox. Keyed by mailbox like every other per-mailbox event, so one
// mailbox's relays stay in order relative to each other.
func (p *publisher) PublishMessageSeen(ctx context.Context, workerID uuid.UUID, action *models.MessageSeenAction) error {
workerEvent := models.WorkerEvent{
Type: models.WorkerEventTypeMessageSeen,
Body: action,
}
workerTopic := kafka.GetWorkerTopic(workerID.String())
return p.publish(workerTopic, action.EmailID.String(), workerEvent)
}
// PublishAddEmail publishes an add email event to the worker
func (p *publisher) PublishAddEmail(ctx context.Context, workerID uuid.UUID, email *models.AddWorkerEmail) error {
workerEvent := models.WorkerEvent{
+3
View File
@@ -8,6 +8,9 @@ const (
WorkerEventTypeRemoveEmail WorkerEventType = "REMOVE_EMAIL"
WorkerEventTypeEmailValidation WorkerEventType = "EMAIL_VALIDATION"
WorkerEventTypeWarmupAction WorkerEventType = "WARMUP_ACTION"
// WorkerEventTypeMessageSeen relays a read/unread change a person made in
// the unibox out to the mailbox provider, so the two agree.
WorkerEventTypeMessageSeen WorkerEventType = "MESSAGE_SEEN"
)
type WorkerEvent struct {
+49
View File
@@ -448,3 +448,52 @@ type UniboxScheduledItem struct {
// Thread the reply will land in (when the user queued from unibox).
ThreadID *string `json:"thread_id,omitempty"`
}
// MessageSeenAction relays a read/unread change from the unibox to the
// mailbox's provider, for one mailbox and up to SeenRelayChunk messages.
//
// The unibox used to be the only place that knew: a thread read here stayed
// bold in Gmail and one archived here stayed in the inbox, which reads as a
// broken client rather than a design decision. Only the explicit change
// travels; nothing reconciles the two stores in the background, because the
// provider's own state is what the next sync brings back anyway.
type MessageSeenAction struct {
// EmailID is the mailbox, which is how the worker finds the live client.
EmailID uuid.UUID `json:"email_id"`
// Seen is the state to apply to every message in the batch.
Seen bool `json:"seen"`
Messages []MessageSeenRef `json:"messages"`
}
// MessageSeenRef names one message in whichever way its provider needs.
// Every field is optional because the three providers use different halves:
// Gmail and Graph a message id, IMAP a folder and a UID.
type MessageSeenRef struct {
// ProviderID is the Gmail message id or the Graph message id. The column
// behind it is provider-agnostic despite its name.
ProviderID string `json:"provider_id,omitempty"`
UID uint32 `json:"uid,omitempty"`
// Folder is the IMAP folder holding UID. UIDs are only unique within one.
Folder string `json:"folder,omitempty"`
// RFCMessageID is the immutable Message-ID. Graph ids change when a
// message moves, so the worker re-resolves from this when it is present.
RFCMessageID string `json:"rfc_message_id,omitempty"`
}
// SeenRelayChunk bounds one MESSAGE_SEEN event. Gmail accepts 1000 ids per
// batchModify; this leaves room under it and keeps one "mark all as read" on
// a large folder from becoming a single enormous event on the bus.
const SeenRelayChunk = 500
// SeenRelayTarget is one message resolved for the relay: which mailbox it
// belongs to, which worker holds that mailbox, and how its provider names it.
type SeenRelayTarget struct {
EmailID uuid.UUID
WorkerID uuid.UUID
// Seen is the state the row holds now, read back rather than taken from
// the request that caused the relay. Two people toggling the same
// conversation in opposite directions at once would otherwise be able to
// leave the provider holding the earlier answer.
Seen bool
Ref MessageSeenRef
}
+90 -17
View File
@@ -43,13 +43,19 @@ type UniboxRepository interface {
Search(ctx context.Context, orgID uuid.UUID, params *models.MailSearchParams) (*models.MailSearchResult, error)
GetUnseenCount(ctx context.Context, orgID uuid.UUID, emailAccountID *uuid.UUID) (int64, error)
MarkSeen(ctx context.Context, userID, id uuid.UUID, seen bool) error
MarkSeenBulk(ctx context.Context, orgID uuid.UUID, ids []uuid.UUID, seen bool) error
// MarkSeenBulk flips the read state of the given messages and returns the
// ids that actually changed, which is what gets relayed to the provider.
MarkSeenBulk(ctx context.Context, orgID uuid.UUID, ids []uuid.UUID, seen bool) ([]uuid.UUID, error)
// MarkSeenByFolder flips the read state of every message in one canonical
// folder for the whole workspace (the sidebar's "mark all as read").
MarkSeenByFolder(ctx context.Context, orgID uuid.UUID, folder string, seen bool) error
MarkSeenByFolder(ctx context.Context, orgID uuid.UUID, folder string, seen bool) ([]uuid.UUID, error)
// MoveToFolderBulk re-files the given messages into one canonical folder,
// org-scoped like MarkSeenBulk.
MoveToFolderBulk(ctx context.Context, orgID uuid.UUID, ids []uuid.UUID, folder string) error
// SeenRelayTargets names the given messages the way their provider does,
// with the worker holding each mailbox. Rows whose mailbox has no worker
// are left out: there is nothing to relay through.
SeenRelayTargets(ctx context.Context, orgID uuid.UUID, ids []uuid.UUID) ([]models.SeenRelayTarget, error)
Delete(ctx context.Context, userID, id uuid.UUID) error
// Snooze: per (user, thread). UpsertSnooze adopts the new
@@ -295,6 +301,9 @@ func (r *uniboxRepository) GetByID(ctx context.Context, userID, id uuid.UUID) (*
// user_id: the body's object-storage key is built from the owner (and the
// row's email_id), so the caller must fetch the body under the owner, not
// under itself.
//
// It does not mark the message read: that is a write with a provider relay
// behind it, and it belongs to the service.
func (r *uniboxRepository) GetByIDForOrg(ctx context.Context, orgID, id uuid.UUID) (*models.EmailMessageStoreData, uuid.UUID, error) {
query := fmt.Sprintf(`
SELECT user_id, %s
@@ -319,12 +328,9 @@ func (r *uniboxRepository) GetByIDForOrg(ctx context.Context, orgID, id uuid.UUI
return nil, uuid.Nil, err
}
// Auto-mark as seen, org-scoped so any member clears the shared unread state.
if !e.Seen {
_ = r.MarkSeenBulk(ctx, orgID, []uuid.UUID{id}, true)
e.Seen = true
}
// Reading it is what marks it read, and that now has to reach the mailbox
// too, so the service owns the transition (it holds the relay). The row
// is returned exactly as stored.
return &e, ownerID, nil
}
@@ -656,32 +662,99 @@ func (r *uniboxRepository) MarkSeen(ctx context.Context, userID, id uuid.UUID, s
return err
}
func (r *uniboxRepository) MarkSeenBulk(ctx context.Context, orgID uuid.UUID, ids []uuid.UUID, seen bool) error {
func (r *uniboxRepository) MarkSeenBulk(ctx context.Context, orgID uuid.UUID, ids []uuid.UUID, seen bool) ([]uuid.UUID, error) {
if len(ids) == 0 {
return nil
return nil, nil
}
// Org-scoped so any member with unibox access can clear the shared inbox's
// unread state, not only the mailbox owner. The unread count is org-wide, so
// a user_id filter would leave the badge stuck for non-owner members. ANY($3)
// also covers the single-id case.
_, err := r.db.Exec(ctx,
//
// `seen <> $1` and the RETURNING are what keep the provider relay honest:
// re-reading a thread that is already read should cost nothing at Gmail.
rows, err := r.db.Query(ctx,
`UPDATE unibox_emails SET seen = $1, updated_at = NOW()
WHERE id = ANY($3) AND email_id IN (SELECT id FROM email_accounts WHERE organization_id = $2)`,
WHERE id = ANY($3) AND seen <> $1
AND email_id IN (SELECT id FROM email_accounts WHERE organization_id = $2)
RETURNING id`,
seen, orgID, ids,
)
return err
if err != nil {
return nil, err
}
defer rows.Close()
changed := make([]uuid.UUID, 0, len(ids))
for rows.Next() {
var id uuid.UUID
if err := rows.Scan(&id); err != nil {
return nil, err
}
changed = append(changed, id)
}
return changed, rows.Err()
}
// MarkSeenByFolder flips the read state of every message in one folder,
// org-scoped like MarkSeenBulk (the sidebar's "mark all as read").
func (r *uniboxRepository) MarkSeenByFolder(ctx context.Context, orgID uuid.UUID, folder string, seen bool) error {
_, err := r.db.Exec(ctx,
func (r *uniboxRepository) MarkSeenByFolder(ctx context.Context, orgID uuid.UUID, folder string, seen bool) ([]uuid.UUID, error) {
rows, err := r.db.Query(ctx,
`UPDATE unibox_emails SET seen = $1, updated_at = NOW()
WHERE folder = $3 AND seen <> $1
AND email_id IN (SELECT id FROM email_accounts WHERE organization_id = $2)`,
AND email_id IN (SELECT id FROM email_accounts WHERE organization_id = $2)
RETURNING id`,
seen, orgID, folder,
)
return err
if err != nil {
return nil, err
}
defer rows.Close()
var changed []uuid.UUID
for rows.Next() {
var id uuid.UUID
if err := rows.Scan(&id); err != nil {
return nil, err
}
changed = append(changed, id)
}
return changed, rows.Err()
}
// SeenRelayTargets resolves messages to what their provider calls them, plus
// the worker holding the mailbox.
//
// The three providers need different halves of this row (Gmail and Graph a
// message id, IMAP a folder and a UID), so all of it travels and the worker
// takes what its client uses. The read state comes from the row rather than
// from the request, so what is relayed is what Warmbly currently holds.
func (r *uniboxRepository) SeenRelayTargets(ctx context.Context, orgID uuid.UUID, ids []uuid.UUID) ([]models.SeenRelayTarget, error) {
if len(ids) == 0 {
return nil, nil
}
rows, err := r.db.Query(ctx,
`SELECT ue.email_id, ea.worker_id, ue.seen, ue.gmail_id, ue.uid, ue.folder_path, ue.message_id
FROM unibox_emails ue
JOIN email_accounts ea ON ea.id = ue.email_id
WHERE ue.id = ANY($2) AND ea.organization_id = $1 AND ea.worker_id IS NOT NULL
ORDER BY ue.email_id`,
orgID, ids,
)
if err != nil {
return nil, err
}
defer rows.Close()
var out []models.SeenRelayTarget
for rows.Next() {
var t models.SeenRelayTarget
if err := rows.Scan(&t.EmailID, &t.WorkerID, &t.Seen, &t.Ref.ProviderID, &t.Ref.UID, &t.Ref.Folder, &t.Ref.RFCMessageID); err != nil {
return nil, err
}
out = append(out, t)
}
return out, rows.Err()
}
func (r *uniboxRepository) MoveToFolderBulk(ctx context.Context, orgID uuid.UUID, ids []uuid.UUID, folder string) error {