Merge pull request #793 from warmbly/fix/inbox-follow-up-sweep-paging

feat: page the hourly inbox follow-up sweep through each mailbox newest first in bounded keyset pages that resolve thread state only for the threads on the page, resume each workspace's cycle from a cursor persisted in inbox_follow_up_sweeps under a per-pass time budget, check recently active threads first every pass, and add the unibox_emails (email_id, internal_date, id) index the walk reads
This commit is contained in:
Matthew Meszaros
2026-10-02 06:33:14 +00:00
committed by GitHub
22 changed files with 1865 additions and 180 deletions
+3 -2
View File
@@ -185,7 +185,7 @@ func runInboxTagFollowUps(ctx context.Context, args []string) error {
fs := newFlagSet("inbox-tag follow-ups")
org := fs.String("org", "", "organization id, or the owner's email address")
days := fs.Int("days", 90, "how far back to consider threads")
limit := fs.Int("limit", 2000, "most threads to sweep in this run")
_ = fs.Int("limit", 0, "deprecated and ignored: the sweep now checks every thread")
if err := fs.Parse(args); err != nil {
return err
}
@@ -215,7 +215,8 @@ func runInboxTagFollowUps(ctx context.Context, args []string) error {
true,
)
p, err := svc.SweepFollowUps(ctx, orgID, time.Now().AddDate(0, 0, -*days), *limit)
// A full cycle of its own, so the hourly sweep's position is left where it was.
p, err := svc.SweepFollowUps(ctx, orgID, inboxtag.FollowUpSweep{Since: time.Now().AddDate(0, 0, -*days), Full: true})
if err != nil {
return fmt.Errorf("follow-up sweep: %w", err)
}
+3 -1
View File
@@ -87,7 +87,9 @@ A thread you wrote to recently wears nothing: there is nothing to do yet, and th
These labels change as the calendar moves, so they are recomputed hourly and **replaced rather than added**. A thread is never wearing both **Follow up** and **Needs reply**; it wears its state, not its history. A reply clears the follow-up label within the hour.
To recompute on demand:
Every conversation with activity in the last 90 days is checked, however many there are. A large workspace is checked a page at a time: each hourly pass spends a few minutes at most and picks up where the previous one stopped, including after a restart. An older conversation therefore gets its label within one full sweep of the workspace after its day comes, which is one pass for most workspaces and several for a very large one. Each pass first checks every conversation that received mail or had a reply classified since the previous pass, whatever date the message carries, so a reply that synced late still clears its label on the next pass.
To recompute on demand (this checks every conversation in one run and leaves the hourly sweep's place alone):
```bash
warmblyctl inbox-tag follow-ups --org you@example.com
+12 -3
View File
@@ -135,6 +135,12 @@ const followUpSweepInterval = time.Hour
// touched in three months is not one anybody is about to chase.
const followUpSweepWindow = 90
// followUpSweepBudget is how long one pass may page through a workspace before the next resumes it.
const followUpSweepBudget = 2 * time.Minute
// followUpSweepFresh caps how far back a pass checks threads changed since the last pass; older changes wait for the cycle.
const followUpSweepFresh = 24 * time.Hour
func (s *JobsService) sweepFollowUps(ctx context.Context) {
if s.InboxTagger == nil || s.EmailRepository == nil {
return
@@ -148,14 +154,17 @@ func (s *JobsService) sweepFollowUps(ctx context.Context) {
log.Warn().Err(err).Msg("follow-up sweep: could not list workspaces")
}
for _, orgID := range orgs {
since := time.Now().AddDate(0, 0, -followUpSweepWindow)
if p, serr := s.InboxTagger.SweepFollowUps(ctx, orgID, since, 0); serr != nil {
if p, serr := s.InboxTagger.SweepFollowUps(ctx, orgID, inboxtag.FollowUpSweep{
Since: time.Now().AddDate(0, 0, -followUpSweepWindow),
Fresh: followUpSweepFresh,
Budget: followUpSweepBudget,
}); serr != nil {
if ctx.Err() != nil {
return
}
log.Warn().Err(serr).Str("org_id", orgID.String()).Msg("follow-up sweep failed")
} else if p.Threads > 0 {
log.Debug().Str("org_id", orgID.String()).Int("threads", p.Threads).Msg("follow-up sweep")
log.Debug().Str("org_id", orgID.String()).Int("threads", p.Threads).Int("pages", p.Pages).Int("skipped", p.Skipped).Bool("cycle_complete", p.Complete).Msg("follow-up sweep")
if s.StreamingPublisher != nil {
s.StreamingPublisher.PublishEmailUpdated(ctx, &pubsub.EmailInboxEvent{OrgID: orgID.String()})
}
+1 -1
View File
@@ -338,7 +338,7 @@ func TestSweepSeedsWorkspaceLabels(t *testing.T) {
svc.WireSettings(fakeSettings{s: models.InboxTaggingSettings{
Questions: []models.InboxTagQuestion{laterMaybe(models.InboxTagQuestionAction{}), roleQuestion()},
}})
if _, err := svc.SweepFollowUps(context.Background(), uuid.New(), time.Now().Add(-time.Hour), 10); err != nil {
if _, err := svc.SweepFollowUps(context.Background(), uuid.New(), FollowUpSweep{Since: time.Now().Add(-time.Hour)}); err != nil {
t.Fatalf("sweep: %v", err)
}
for _, want := range []string{"later-maybe", "recruiter", "hiring-manager"} {
+407 -4
View File
@@ -1,11 +1,17 @@
package inboxtag
import (
"bytes"
"context"
"errors"
"fmt"
"strings"
"testing"
"time"
"github.com/google/uuid"
"github.com/rs/zerolog"
"github.com/rs/zerolog/log"
"github.com/warmbly/warmbly/internal/repository"
)
@@ -138,6 +144,10 @@ type fakeCategories struct {
labels map[string]map[string]bool
seeded []string
removed map[string][]string
// synced counts follow-up syncs per thread; onSync runs after each.
synced map[string]int
onSync func()
fail map[string]error
}
func (f *fakeCategories) EnsureCategory(_ context.Context, _ uuid.UUID, slug string) (uuid.UUID, error) {
@@ -160,6 +170,9 @@ func (f *fakeCategories) AddThreadLabels(_ context.Context, _ uuid.UUID, threadI
return nil
}
func (f *fakeCategories) SyncExclusiveLabels(ctx context.Context, orgID uuid.UUID, threadID string, family []string, want string) error {
if err := f.fail[threadID]; err != nil {
return err
}
if f.labels == nil {
f.labels = map[string]map[string]bool{}
}
@@ -174,6 +187,13 @@ func (f *fakeCategories) SyncExclusiveLabels(ctx context.Context, orgID uuid.UUI
if want != "" {
f.labels[threadID][want] = true
}
if f.synced == nil {
f.synced = map[string]int{}
}
f.synced[threadID]++
if f.onSync != nil {
f.onSync()
}
return nil
}
func (f *fakeCategories) RemoveAutoLabels(_ context.Context, _ uuid.UUID, threadID string, slugs []string) error {
@@ -198,7 +218,7 @@ func TestSweepNeedsNoModel(t *testing.T) {
asker := &countingAsker{}
svc := NewService(asker, repo, cats, nil, true)
p, err := svc.SweepFollowUps(context.Background(), uuid.New(), now.AddDate(0, 0, -90), 0)
p, err := svc.SweepFollowUps(context.Background(), uuid.New(), FollowUpSweep{Since: now.AddDate(0, 0, -90)})
if err != nil {
t.Fatalf("sweep: %v", err)
}
@@ -235,7 +255,7 @@ func TestSweepReplacesRatherThanAccumulates(t *testing.T) {
svc := NewService(&countingAsker{}, repo, cats, nil, true)
orgID := uuid.New()
if _, err := svc.SweepFollowUps(context.Background(), orgID, now.AddDate(0, 0, -90), 0); err != nil {
if _, err := svc.SweepFollowUps(context.Background(), orgID, FollowUpSweep{Since: now.AddDate(0, 0, -90)}); err != nil {
t.Fatalf("first sweep: %v", err)
}
if !cats.has("t-1", LabelFollowUp) {
@@ -246,7 +266,7 @@ func TestSweepReplacesRatherThanAccumulates(t *testing.T) {
repo.states[0].LastInboundAt = now.AddDate(0, 0, -3)
repo.states[0].LastKind = KindHumanReply
repo.states[0].BestIntent = IntentWantsInfo
if _, err := svc.SweepFollowUps(context.Background(), orgID, now.AddDate(0, 0, -90), 0); err != nil {
if _, err := svc.SweepFollowUps(context.Background(), orgID, FollowUpSweep{Since: now.AddDate(0, 0, -90)}); err != nil {
t.Fatalf("second sweep: %v", err)
}
if !cats.has("t-1", LabelNeedsReply) {
@@ -258,10 +278,393 @@ func TestSweepReplacesRatherThanAccumulates(t *testing.T) {
// We answer, so nothing is owed either way yet.
repo.states[0].LastOutboundAt = now
if _, err := svc.SweepFollowUps(context.Background(), orgID, now.AddDate(0, 0, -90), 0); err != nil {
if _, err := svc.SweepFollowUps(context.Background(), orgID, FollowUpSweep{Since: now.AddDate(0, 0, -90)}); err != nil {
t.Fatalf("third sweep: %v", err)
}
if cats.has("t-1", LabelNeedsReply) {
t.Error("still asking us to reply after we did")
}
}
// quietThreads is n threads we wrote to an hour ago, with one we wrote to
// eight days ago at index stale, which is due a Follow up.
func quietThreads(n, stale int) []repository.ThreadFollowUpState {
now := time.Now()
out := make([]repository.ThreadFollowUpState, n)
for i := range out {
out[i] = repository.ThreadFollowUpState{ThreadID: fmt.Sprintf("t-%d", i), LastOutboundAt: now.Add(-time.Hour)}
}
out[stale].LastOutboundAt = now.AddDate(0, 0, -8)
return out
}
// A thread outside the newest 2,000 is reached: each pass takes one page and
// the next resumes after it, until the cycle ends and starts again at the newest.
func TestSweepPagesThroughEveryThread(t *testing.T) {
repo := &fakeRepo{states: quietThreads(2500, 2400)}
cats := &fakeCategories{}
svc := NewService(&countingAsker{}, repo, cats, nil, true)
orgID := uuid.New()
opts := FollowUpSweep{Since: time.Now().AddDate(0, 0, -90), Budget: time.Nanosecond, PageSize: 500}
passes := 0
for {
passes++
p, err := svc.SweepFollowUps(context.Background(), orgID, opts)
if err != nil {
t.Fatalf("pass %d: %v", passes, err)
}
if p.Pages != 1 {
t.Fatalf("pass %d read %d pages under a spent budget, want 1", passes, p.Pages)
}
if p.Complete {
break
}
if repo.cursor == nil {
t.Fatalf("pass %d stopped mid-cycle without saving where", passes)
}
if passes > 10 {
t.Fatal("the cycle never completed")
}
}
if passes != 6 {
t.Errorf("took %d passes, want 6 (five full pages and the empty one that ends the cycle)", passes)
}
if !cats.has("t-2400", LabelFollowUp) {
t.Error("the stale thread beyond the newest 2,000 was never labelled")
}
for i := range 2500 {
if n := cats.synced[fmt.Sprintf("t-%d", i)]; n != 1 {
t.Fatalf("t-%d was evaluated %d times in one cycle, want once", i, n)
}
}
if repo.cursor != nil {
t.Fatalf("a finished cycle left a cursor at %+v", repo.cursor)
}
// The next cycle starts again at the newest thread.
if _, err := svc.SweepFollowUps(context.Background(), orgID, opts); err != nil {
t.Fatalf("next cycle: %v", err)
}
if cats.synced["t-0"] != 2 || cats.synced["t-500"] != 1 {
t.Errorf("next cycle evaluated t-0 %d and t-500 %d times, want it to restart at the newest page", cats.synced["t-0"], cats.synced["t-500"])
}
}
// A pass cut off mid-page keeps the cursor on the last thread it evaluated, and
// the next pass picks up right after it.
func TestSweepResumesAfterTheLastEvaluatedThread(t *testing.T) {
repo := &fakeRepo{states: quietThreads(1200, 1100)}
cats := &fakeCategories{}
svc := NewService(&countingAsker{}, repo, cats, nil, true)
orgID := uuid.New()
opts := FollowUpSweep{Since: time.Now().AddDate(0, 0, -90), PageSize: 500}
ctx, cancel := context.WithCancel(context.Background())
calls := 0
cats.onSync = func() {
if calls++; calls == 700 {
cancel()
}
}
if _, err := svc.SweepFollowUps(ctx, orgID, opts); !errors.Is(err, context.Canceled) {
t.Fatalf("interrupted sweep returned %v, want context.Canceled", err)
}
if repo.cursor == nil || repo.cursor.RowID != repo.position(699).RowID {
t.Fatalf("cursor %+v, want the 700th thread, the last one evaluated", repo.cursor)
}
cats.onSync = nil
p, err := svc.SweepFollowUps(context.Background(), orgID, opts)
if err != nil {
t.Fatalf("resumed sweep: %v", err)
}
if p.Threads != 500 || !p.Complete {
t.Fatalf("resumed sweep evaluated %d threads (complete %v), want the 500 left", p.Threads, p.Complete)
}
for i := range 1200 {
if n := cats.synced[fmt.Sprintf("t-%d", i)]; n != 1 {
t.Fatalf("t-%d evaluated %d times across the interrupted and resumed passes, want once", i, n)
}
}
if !cats.has("t-1100", LabelFollowUp) {
t.Error("the stale thread after the interruption was never labelled")
}
}
// A reply stored since the last pass is checked in the next one whatever date it
// carries, while the cycle is elsewhere.
func TestSweepChecksChangedThreadsEveryPass(t *testing.T) {
now := time.Now()
repo := &fakeRepo{states: quietThreads(3000, 2999), base: now}
// They answered a day ago in a thread we chased a week ago; the sync stored it ten minutes ago.
repo.states[2500] = repository.ThreadFollowUpState{ThreadID: "t-2500", LastOutboundAt: now.AddDate(0, 0, -7), LastInboundAt: now.AddDate(0, 0, -1), BestIntent: IntentWantsInfo, LastKind: KindHumanReply}
repo.changes = []repository.FollowUpChange{
{At: now.Add(-2 * time.Hour), RowID: uuid.New(), ThreadID: "t-2998"},
{At: now.Add(-10 * time.Minute), RowID: uuid.New(), ThreadID: "t-2500"},
}
repo.fresh = &repository.FollowUpMark{At: now.Add(-time.Hour)}
pos := repo.position(1000)
repo.cursor = &pos
cats := &fakeCategories{labels: map[string]map[string]bool{"t-2500": {LabelFollowUp: true}}}
svc := NewService(&countingAsker{}, repo, cats, nil, true)
p, err := svc.SweepFollowUps(context.Background(), uuid.New(), FollowUpSweep{
Since: now.AddDate(0, 0, -90), Fresh: 24 * time.Hour, Budget: time.Nanosecond, PageSize: 500,
})
if err != nil {
t.Fatalf("sweep: %v", err)
}
if cats.has("t-2500", LabelFollowUp) {
t.Error("a thread they answered still wears Follow up while the cycle is elsewhere")
}
if cats.synced["t-2998"] != 0 {
t.Error("a change from before the last check was read again")
}
if cats.synced["t-1001"] != 1 {
t.Error("the cycle did not move on from its cursor in the same pass")
}
if p.Threads != 1+500 {
t.Errorf("swept %d threads, want the changed one and one cycle page", p.Threads)
}
if repo.fresh == nil || !repo.fresh.At.After(now.Add(-5*time.Minute)) {
t.Errorf("the changed-thread mark stayed at %+v", repo.fresh)
}
}
// A failing changed-thread check costs that check, not the cycle.
func TestSweepCycleRunsWhenTheChangedThreadCheckFails(t *testing.T) {
repo := &fakeRepo{states: quietThreads(600, 10), failChanges: errors.New("statement timeout")}
cats := &fakeCategories{}
svc := NewService(&countingAsker{}, repo, cats, nil, true)
p, err := svc.SweepFollowUps(context.Background(), uuid.New(), FollowUpSweep{Since: time.Now().AddDate(0, 0, -90), Fresh: time.Hour, PageSize: 500})
if err != nil {
t.Fatalf("sweep: %v", err)
}
if !p.Complete || !cats.has("t-10", LabelFollowUp) {
t.Fatalf("the cycle did not run after the check failed: %+v", p)
}
if repo.freshFailures != 1 {
t.Errorf("fresh failures %d, want the one counted", repo.freshFailures)
}
}
// failAgain runs n passes that must each fail without stepping past anything.
func failAgain(t *testing.T, svc *Service, opts FollowUpSweep, n int, check func(pass int)) {
t.Helper()
for pass := 1; pass <= n; pass++ {
if p, err := svc.SweepFollowUps(context.Background(), uuid.New(), opts); err == nil || p.Skipped != 0 {
t.Fatalf("pass %d: %+v %v, want the failure reported and nothing skipped", pass, p, err)
}
check(pass)
}
}
// backdate makes the current run of failures look an hour and a half old.
func backdate(since **time.Time) {
old := time.Now().Add(-90 * time.Minute)
*since = &old
}
// A page that fails is retried, and stepped past only once it has failed
// repeatedly over at least an hour, so the cycle keeps moving.
func TestSweepStepsPastAPageThatKeepsFailing(t *testing.T) {
repo := &fakeRepo{states: quietThreads(1200, 1100)}
poison := repo.position(499).RowID
repo.failPage = func(_ uuid.UUID, after *repository.FollowUpPosition) error {
if after != nil && after.RowID == poison {
return errors.New("statement timeout")
}
return nil
}
cats := &fakeCategories{}
svc := NewService(&countingAsker{}, repo, cats, nil, true)
opts := FollowUpSweep{Since: time.Now().AddDate(0, 0, -90), PageSize: 500}
// Back to back, as several consumers would run them during a short outage.
failAgain(t, svc, opts, followUpFailureLimit+1, func(pass int) {
if repo.cursor == nil || repo.cursor.RowID != poison || repo.pageFailures != pass || repo.pageFailingSince == nil {
t.Fatalf("pass %d left cursor %+v with %d failures since %v", pass, repo.cursor, repo.pageFailures, repo.pageFailingSince)
}
})
backdate(&repo.pageFailingSince)
p, err := svc.SweepFollowUps(context.Background(), uuid.New(), opts)
if err != nil || !p.Complete || p.Skipped != 1 {
t.Fatalf("after an hour of failures: %+v %v, want the page skipped and the cycle finished", p, err)
}
if !cats.has("t-1100", LabelFollowUp) {
t.Error("the thread past the failing page was never labelled")
}
if cats.synced["t-700"] != 0 || repo.pageFailures != 0 || repo.pageFailingSince != nil {
t.Errorf("t-700 synced %d times and %d failures left, want the page skipped and the count reset", cats.synced["t-700"], repo.pageFailures)
}
}
// A label that cannot be written stops the cursor before that thread, and
// after an hour of failures the thread alone is stepped past.
func TestSweepStopsAtAThreadWhoseLabelFails(t *testing.T) {
repo := &fakeRepo{states: quietThreads(1200, 1100)}
cats := &fakeCategories{fail: map[string]error{"t-300": errors.New("deadlock detected")}}
svc := NewService(&countingAsker{}, repo, cats, nil, true)
opts := FollowUpSweep{Since: time.Now().AddDate(0, 0, -90), PageSize: 500}
failAgain(t, svc, opts, followUpFailureLimit, func(pass int) {
if repo.cursor == nil || repo.cursor.RowID != repo.position(299).RowID || repo.pageFailures != pass {
t.Fatalf("pass %d left the cursor at %+v with %d failures, want the thread before the failure", pass, repo.cursor, repo.pageFailures)
}
})
backdate(&repo.pageFailingSince)
p, err := svc.SweepFollowUps(context.Background(), uuid.New(), opts)
if err != nil || !p.Complete || p.Skipped != 1 {
t.Fatalf("after an hour of failures: %+v %v, want the thread skipped and the cycle finished", p, err)
}
if cats.synced["t-299"] != 1 || cats.synced["t-301"] != 1 || !cats.has("t-1100", LabelFollowUp) {
t.Errorf("t-299 %d, t-301 %d syncs, t-1100 labelled %v", cats.synced["t-299"], cats.synced["t-301"], cats.has("t-1100", LabelFollowUp))
}
}
// When even the positions of a failing page cannot be read, the rest of that
// mailbox is left to the next cycle and the walk moves to the next mailbox.
func TestSweepStepsPastAMailboxItCannotRead(t *testing.T) {
repo := &fakeRepo{
states: quietThreads(600, 10),
second: []repository.ThreadFollowUpState{{ThreadID: "m2-stale", LastOutboundAt: time.Now().AddDate(0, 0, -8)}},
failPositions: errors.New("statement timeout"),
}
repo.failPage = func(mailbox uuid.UUID, after *repository.FollowUpPosition) error {
if mailbox == fakeMailbox && after != nil {
return errors.New("statement timeout")
}
return nil
}
cats := &fakeCategories{}
svc := NewService(&countingAsker{}, repo, cats, nil, true)
opts := FollowUpSweep{Since: time.Now().AddDate(0, 0, -90), PageSize: 500}
failAgain(t, svc, opts, followUpFailureLimit, func(int) {})
backdate(&repo.pageFailingSince)
p, err := svc.SweepFollowUps(context.Background(), uuid.New(), opts)
if err != nil || !p.Complete || p.Skipped != 1 {
t.Fatalf("after an hour of failures: %+v %v, want the mailbox skipped and the cycle finished", p, err)
}
if !cats.has("m2-stale", LabelFollowUp) {
t.Error("the next mailbox was never reached")
}
}
// A failed read of changed threads never moves the mark, however long it has been failing.
func TestSweepNeverMovesTheMarkPastAFailedRead(t *testing.T) {
mark := repository.FollowUpMark{At: time.Now().Add(-time.Hour), RowID: uuid.New()}
repo := &fakeRepo{states: quietThreads(10, 5), fresh: &mark, freshFailures: 5, failChanges: errors.New("connection reset")}
backdate(&repo.freshFailingSince)
svc := NewService(&countingAsker{}, repo, &fakeCategories{}, nil, true)
for pass := 1; pass <= 2; pass++ {
p, err := svc.SweepFollowUps(context.Background(), uuid.New(), FollowUpSweep{Since: time.Now().AddDate(0, 0, -90), Fresh: 24 * time.Hour, PageSize: 500})
if err != nil || !p.Complete || p.Skipped != 0 {
t.Fatalf("pass %d: %+v %v, want the cycle run and nothing skipped", pass, p, err)
}
if repo.fresh == nil || *repo.fresh != mark {
t.Fatalf("pass %d moved the mark to %+v past changes it never read", pass, repo.fresh)
}
}
if repo.freshFailures != 7 {
t.Errorf("fresh failures %d, want each failed read counted", repo.freshFailures)
}
}
// A pass that recovers where it last failed does not carry the old count forward.
func TestSweepClearsFailuresOnceItMovesOn(t *testing.T) {
repo := &fakeRepo{states: quietThreads(1200, 1100), pageFailures: 2}
backdate(&repo.pageFailingSince)
cats := &fakeCategories{}
ctx, cancel := context.WithCancel(context.Background())
calls := 0
cats.onSync = func() {
if calls++; calls == 10 {
cancel()
}
}
svc := NewService(&countingAsker{}, repo, cats, nil, true)
if _, err := svc.SweepFollowUps(ctx, uuid.New(), FollowUpSweep{Since: time.Now().AddDate(0, 0, -90), PageSize: 500}); !errors.Is(err, context.Canceled) {
t.Fatalf("interrupted sweep returned %v", err)
}
if repo.cursor == nil || repo.cursor.RowID != repo.position(9).RowID {
t.Fatalf("cursor %+v, want the tenth thread", repo.cursor)
}
if repo.pageFailures != 0 || repo.pageFailingSince != nil {
t.Errorf("saved %d failures since %v after moving on", repo.pageFailures, repo.pageFailingSince)
}
}
// A pass whose lease another walker takes over stops, says so, and reports Busy.
func TestSweepStopsWhenItsLeaseIsTakenOver(t *testing.T) {
var buf bytes.Buffer
prev := log.Logger
log.Logger = zerolog.New(&buf)
t.Cleanup(func() { log.Logger = prev })
repo := &fakeRepo{states: quietThreads(1200, 1100)}
cats := &fakeCategories{}
calls := 0
cats.onSync = func() {
switch calls++; calls {
case 10:
// A slow pass outlives its lease; nobody has taken over, so it may still save.
repo.leasedUntil = time.Now().Add(-time.Minute)
case 700:
repo.leaseOwner = uuid.New()
}
}
svc := NewService(&countingAsker{}, repo, cats, nil, true)
p, err := svc.SweepFollowUps(context.Background(), uuid.New(), FollowUpSweep{Since: time.Now().AddDate(0, 0, -90), PageSize: 500})
if err != nil || !p.Busy {
t.Fatalf("taken-over pass: %+v %v, want Busy and no error", p, err)
}
if repo.cursor == nil || repo.cursor.RowID != repo.position(499).RowID {
t.Errorf("cursor %+v, want the first page saved after the lease lapsed", repo.cursor)
}
if !strings.Contains(buf.String(), "taken over") {
t.Errorf("a lost lease was not logged: %q", buf.String())
}
}
// A second walker finds the workspace leased and leaves its state alone.
func TestSweepLeaseKeepsOneWalker(t *testing.T) {
repo := &fakeRepo{states: quietThreads(1200, 1100)}
pos := repo.position(600)
repo.cursor = &pos
holder := uuid.New()
repo.leaseOwner, repo.leasedUntil = holder, time.Now().Add(time.Minute)
cats := &fakeCategories{}
svc := NewService(&countingAsker{}, repo, cats, nil, true)
p, err := svc.SweepFollowUps(context.Background(), uuid.New(), FollowUpSweep{Since: time.Now().AddDate(0, 0, -90), PageSize: 500})
if err != nil || !p.Busy || p.Threads != 0 || repo.pages != 0 {
t.Fatalf("second walker: %+v %v after %d pages, want it to stand aside", p, err, repo.pages)
}
if repo.cursor == nil || repo.cursor.RowID != pos.RowID || repo.leaseOwner != holder {
t.Fatalf("second walker moved the cursor to %+v or took the lease", repo.cursor)
}
}
// An operator's full run neither reads nor moves the hourly sweep's cursor.
func TestFullSweepLeavesTheCursorAlone(t *testing.T) {
repo := &fakeRepo{states: quietThreads(1200, 1100)}
pos := repo.position(600)
repo.cursor = &pos
cats := &fakeCategories{}
svc := NewService(&countingAsker{}, repo, cats, nil, true)
p, err := svc.SweepFollowUps(context.Background(), uuid.New(), FollowUpSweep{Since: time.Now().AddDate(0, 0, -90), Full: true})
if err != nil {
t.Fatalf("sweep: %v", err)
}
if p.Threads != 1200 || !p.Complete {
t.Fatalf("full sweep evaluated %d threads (complete %v), want all 1200", p.Threads, p.Complete)
}
if repo.cursor == nil || repo.cursor.RowID != pos.RowID {
t.Fatalf("full sweep moved the cursor to %+v", repo.cursor)
}
}
+144 -2
View File
@@ -1,10 +1,13 @@
package inboxtag
import (
"bytes"
"context"
"encoding/json"
"os"
"path/filepath"
"slices"
"strconv"
"testing"
"time"
@@ -231,6 +234,23 @@ type fakeRepo struct {
inReplyTo []string
reopened []string
states []repository.ThreadFollowUpState
cursor *repository.FollowUpPosition
pages int
base time.Time
// The follow-up sweep's persisted state and lease.
fresh *repository.FollowUpMark
pageFailures int
freshFailures int
leaseOwner uuid.UUID
leasedUntil time.Time
changes []repository.FollowUpChange
failPage func(mailbox uuid.UUID, after *repository.FollowUpPosition) error
failChanges error
failPositions error
second []repository.ThreadFollowUpState
// pageFailingSince and freshFailingSince are when the current run of failures began.
pageFailingSince *time.Time
freshFailingSince *time.Time
}
func (f *fakeRepo) Claim(_ context.Context, _, _ uuid.UUID, id, _ string) (bool, error) {
@@ -281,8 +301,130 @@ func (f *fakeRepo) Reopen(_ context.Context, _ uuid.UUID, id, _ string) ([]strin
return []string{"cold-inbound", "needs-review"}, nil
}
func (f *fakeRepo) ThreadStates(context.Context, uuid.UUID, time.Time, int) ([]repository.ThreadFollowUpState, error) {
return f.states, nil
// fakeMailbox holds the fake threads in states, fakeMailbox2 those in second; each in walk order, one message each.
var (
fakeMailbox = uuid.MustParse("00000000-0000-0000-0000-000000000001")
fakeMailbox2 = uuid.MustParse("00000000-0000-0000-0000-000000000002")
)
// position puts thread i one minute older than thread i-1, counting back from base.
func (f *fakeRepo) position(i int) repository.FollowUpPosition {
return f.positionIn(fakeMailbox, i)
}
func (f *fakeRepo) positionIn(mailbox uuid.UUID, i int) repository.FollowUpPosition {
base := f.base
if base.IsZero() {
base = time.Now()
}
return repository.FollowUpPosition{
MailboxID: mailbox,
At: base.Add(-time.Duration(i) * time.Minute),
RowID: uuid.NewSHA1(mailbox, []byte(strconv.Itoa(i))),
}
}
func (f *fakeRepo) FollowUpMailboxes(context.Context, uuid.UUID) ([]uuid.UUID, error) {
if len(f.second) > 0 {
return []uuid.UUID{fakeMailbox, fakeMailbox2}, nil
}
return []uuid.UUID{fakeMailbox}, nil
}
func (f *fakeRepo) FollowUpPage(_ context.Context, _, mailbox uuid.UUID, since time.Time, after *repository.FollowUpPosition, limit int) (repository.FollowUpPage, error) {
f.pages++
if f.failPage != nil {
if err := f.failPage(mailbox, after); err != nil {
return repository.FollowUpPage{}, err
}
}
return f.page(mailbox, since, after, limit), nil
}
func (f *fakeRepo) page(mailbox uuid.UUID, since time.Time, after *repository.FollowUpPosition, limit int) repository.FollowUpPage {
states := f.states
if mailbox == fakeMailbox2 {
states = f.second
}
start := 0
if after != nil {
for start < len(states) && f.positionIn(mailbox, start).RowID != after.RowID {
start++
}
start++
}
var page repository.FollowUpPage
for i := start; i < len(states) && i < start+limit && !f.positionIn(mailbox, i).At.Before(since); i++ {
st := states[i]
st.Position = f.positionIn(mailbox, i)
page.States = append(page.States, st)
page.Rows++
page.Last = &st.Position
}
return page
}
func (f *fakeRepo) FollowUpPagePositions(_ context.Context, _, mailbox uuid.UUID, since time.Time, after *repository.FollowUpPosition, limit int) (repository.FollowUpPage, error) {
if f.failPositions != nil {
return repository.FollowUpPage{}, f.failPositions
}
page := f.page(mailbox, since, after, limit)
page.States = nil
return page, nil
}
func (f *fakeRepo) FollowUpChanges(_ context.Context, _ uuid.UUID, after repository.FollowUpMark, until time.Time, limit int) ([]repository.FollowUpChange, error) {
if f.failChanges != nil {
return nil, f.failChanges
}
var out []repository.FollowUpChange
for _, c := range f.changes {
later := c.At.After(after.At) || (c.At.Equal(after.At) && bytes.Compare(c.RowID[:], after.RowID[:]) > 0)
if later && !c.At.After(until) && len(out) < limit {
out = append(out, c)
}
}
return out, nil
}
func (f *fakeRepo) FollowUpThreadStates(_ context.Context, _ uuid.UUID, threads []string, _ time.Time) ([]repository.ThreadFollowUpState, error) {
var out []repository.ThreadFollowUpState
for _, st := range append(append([]repository.ThreadFollowUpState{}, f.states...), f.second...) {
if slices.Contains(threads, st.ThreadID) {
out = append(out, st)
}
}
return out, nil
}
func (f *fakeRepo) ClaimFollowUpSweep(_ context.Context, _, owner uuid.UUID, lease time.Duration) (*repository.FollowUpSweepState, error) {
now := time.Now()
if f.leaseOwner != uuid.Nil && f.leaseOwner != owner && f.leasedUntil.After(now) {
return nil, nil
}
f.leaseOwner, f.leasedUntil = owner, now.Add(lease)
return &repository.FollowUpSweepState{
Cursor: f.cursor, Fresh: f.fresh, PageFailures: f.pageFailures, FreshFailures: f.freshFailures,
PageFailingSince: f.pageFailingSince, FreshFailingSince: f.freshFailingSince, Now: now,
}, nil
}
func (f *fakeRepo) SaveFollowUpSweep(_ context.Context, _, owner uuid.UUID, lease time.Duration, st repository.FollowUpSweepState) (bool, error) {
now := time.Now()
if f.leaseOwner != owner {
return false, nil
}
f.cursor, f.fresh, f.pageFailures, f.freshFailures = st.Cursor, st.Fresh, st.PageFailures, st.FreshFailures
f.pageFailingSince, f.freshFailingSince = st.PageFailingSince, st.FreshFailingSince
f.leasedUntil = now.Add(lease)
return true, nil
}
func (f *fakeRepo) ReleaseFollowUpSweep(_ context.Context, _, owner uuid.UUID) error {
if f.leaseOwner == owner {
f.leasedUntil = time.Time{}
}
return nil
}
func (f *fakeRepo) GetByMessageID(_ context.Context, _ uuid.UUID, id string) (*repository.InboxTagResult, error) {
-66
View File
@@ -680,69 +680,3 @@ func (s *Service) PreviousContext(ctx context.Context, accountID uuid.UUID, thre
}
return body, campaign
}
// ── Follow-up sweep ────────────────────────────────────────────────────────
// FollowUpProgress reports what one sweep changed.
type FollowUpProgress struct {
Threads int
Labelled map[string]int
Cleared int
}
// SweepFollowUps recomputes the follow-up label on every recently active thread.
//
// The sweep makes no model calls. It runs when tagging is enabled and reuses
// trusted classifications so automated mail is not treated as a human reply.
func (s *Service) SweepFollowUps(ctx context.Context, orgID uuid.UUID, since time.Time, limit int) (FollowUpProgress, error) {
p := FollowUpProgress{Labelled: map[string]int{}}
if s == nil || s.repo == nil || s.categories == nil {
return p, nil
}
if limit <= 0 {
limit = 2000
}
states, err := s.repo.ThreadStates(ctx, orgID, since, limit)
if err != nil {
return p, err
}
// The hourly sweep is also how a workspace that predates the feature
// gets its labels: the whole taxonomy when classification is on, the
// follow-up labels otherwise. Idempotent and cached, so it costs nothing
// after the first pass.
seed := append([]string{}, SeedSet()...)
if s.Enabled() {
seed = append(seed, CustomLabels(s.workspace(ctx, orgID).questions)...)
}
if err := s.categories.EnsureAll(ctx, orgID, seed); err != nil {
log.Warn().Err(err).Msg("inbox tagging: could not seed labels")
}
now := time.Now()
for _, st := range states {
if err := ctx.Err(); err != nil {
return p, err
}
want := FollowUp(ThreadState{
ThreadID: st.ThreadID,
LastInboundAt: st.LastInboundAt,
LastOutboundAt: st.LastOutboundAt,
BestIntent: st.BestIntent,
LastKind: st.LastKind,
}, now)
if err := s.categories.SyncExclusiveLabels(ctx, orgID, st.ThreadID, FollowUpLabels, want); err != nil {
log.Warn().Err(err).Str("thread_id", st.ThreadID).Msg("inbox tagging: follow-up label not applied")
continue
}
p.Threads++
if want == "" {
p.Cleared++
} else {
p.Labelled[want]++
}
}
return p, nil
}
+445
View File
@@ -0,0 +1,445 @@
package inboxtag
import (
"bytes"
"context"
"errors"
"time"
"github.com/google/uuid"
"github.com/rs/zerolog/log"
"github.com/warmbly/warmbly/internal/repository"
)
const (
// DefaultFollowUpPageSize is how many messages one page of the sweep reads.
DefaultFollowUpPageSize = 500
// followUpFailureLimit and followUpFailureSpan are how many failures, over how long, before a place is stepped past.
followUpFailureLimit = 3
followUpFailureSpan = time.Hour
// followUpSettle keeps the changed-thread check behind writes that may still be committing.
followUpSettle = time.Minute
// followUpLeaseMargin is how long a walker's lease outlives its budget.
followUpLeaseMargin = 2 * time.Minute
// followUpUnboundedLease is the lease of a pass with no budget, renewed on every save.
followUpUnboundedLease = 10 * time.Minute
)
// maxRowID orders after every row id, so a mark at (t, maxRowID) covers everything at t.
var maxRowID = uuid.UUID{0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff, 0xff}
// mailboxDone is a cursor time before any sweep window, marking a mailbox finished for the cycle.
var mailboxDone = time.Unix(0, 0).UTC()
// errLeaseLost stops a walker whose lease another walker has taken over.
var errLeaseLost = errors.New("follow-up sweep: lease lost")
// FollowUpSweep bounds one pass of the follow-up sweep.
type FollowUpSweep struct {
// Since is the oldest activity a thread may have and still be swept.
Since time.Time
// Fresh is the furthest back the check of changed threads looks, its first run included; zero skips it.
Fresh time.Duration
// Budget stops a pass from starting another page; zero runs the cycle to its end.
Budget time.Duration
// PageSize is how many messages a page reads; zero is DefaultFollowUpPageSize.
PageSize int
// Full sweeps every thread from the newest and leaves the scheduled sweep's state alone.
Full bool
}
// FollowUpProgress reports what one sweep changed.
type FollowUpProgress struct {
Threads int
Labelled map[string]int
Cleared int
Pages int
// Skipped counts pages and threads stepped past after repeated failures.
Skipped int
// Complete is a cycle that reached its end, so the next pass starts at the newest.
Complete bool
// Busy is a pass that found another walker holding the workspace.
Busy bool
}
// SweepFollowUps recomputes follow-up labels a page at a time, resuming the workspace's cycle; it makes no model calls.
func (s *Service) SweepFollowUps(ctx context.Context, orgID uuid.UUID, opts FollowUpSweep) (FollowUpProgress, error) {
p := FollowUpProgress{Labelled: map[string]int{}}
if s == nil || s.repo == nil || s.categories == nil {
return p, nil
}
if opts.PageSize <= 0 {
opts.PageSize = DefaultFollowUpPageSize
}
mailboxes, err := s.repo.FollowUpMailboxes(ctx, orgID)
if err != nil {
return p, err
}
// The hourly sweep is also how a workspace that predates the feature
// gets its labels: the whole taxonomy when classification is on, the
// follow-up labels otherwise. Idempotent and cached, so it costs nothing
// after the first pass.
seed := append([]string{}, SeedSet()...)
if s.Enabled() {
seed = append(seed, CustomLabels(s.workspace(ctx, orgID).questions)...)
}
if err := s.categories.EnsureAll(ctx, orgID, seed); err != nil {
log.Warn().Err(err).Msg("inbox tagging: could not seed labels")
}
now := time.Now()
w := &followUpWalk{s: s, orgID: orgID, opts: opts, started: now, now: now, seen: map[string]bool{}, p: &p}
if opts.Full {
complete, err := w.walk(ctx, mailboxes, nil)
p.Complete = complete && err == nil
return p, err
}
w.owner = uuid.New()
w.lease = followUpUnboundedLease
if opts.Budget > 0 {
w.lease = opts.Budget + followUpLeaseMargin
}
st, err := s.repo.ClaimFollowUpSweep(ctx, orgID, w.owner, w.lease)
if err != nil {
return p, err
}
if st == nil {
p.Busy = true
return p, nil
}
w.state = st
defer w.release(ctx)
err = w.run(ctx, mailboxes)
if w.lost {
log.Warn().Str("org_id", orgID.String()).Msg("inbox tagging: follow-up sweep taken over by another walker; this pass stopped")
p.Busy = true
}
if errors.Is(err, errLeaseLost) {
err = nil
}
return p, err
}
type followUpWalk struct {
s *Service
orgID uuid.UUID
opts FollowUpSweep
started time.Time
now time.Time
seen map[string]bool
p *FollowUpProgress
// owner, lease and state are nil-valued on a Full sweep, which persists nothing.
owner uuid.UUID
lease time.Duration
state *repository.FollowUpSweepState
// progressed and freshProgressed are set once the cycle or the changed-thread check moves past where a failure was counted.
progressed bool
freshProgressed bool
// lost is a pass whose lease another walker took over.
lost bool
}
func (w *followUpWalk) run(ctx context.Context, mailboxes []uuid.UUID) error {
if w.opts.Fresh > 0 {
if err := w.fresh(ctx); err != nil {
if errors.Is(err, errLeaseLost) || ctx.Err() != nil {
return err
}
log.Warn().Err(err).Str("org_id", w.orgID.String()).Msg("inbox tagging: changed threads not checked this pass")
}
}
complete, err := w.walk(ctx, mailboxes, w.state.Cursor)
if err != nil || !complete {
return err
}
w.state.Cursor = nil
w.state.PageFailures, w.state.PageFailingSince = 0, nil
if err := w.save(ctx); err != nil {
return err
}
w.p.Complete = true
return nil
}
// fresh checks every thread with a message stored or a verdict written since the last check.
func (w *followUpWalk) fresh(ctx context.Context) error {
st := w.state
mark := repository.FollowUpMark{At: st.Now.Add(-w.opts.Fresh)}
if st.Fresh != nil && st.Fresh.At.After(mark.At) {
mark = *st.Fresh
}
until := st.Now.Add(-followUpSettle)
if !mark.At.Before(until) {
return nil
}
for first := true; ; first = false {
if !first && w.spent() {
return nil
}
changes, err := w.s.repo.FollowUpChanges(ctx, w.orgID, mark, until, w.opts.PageSize)
if err != nil {
// An unread page is never stepped past; the Fresh floor bounds how far the mark can fall behind.
if ctx.Err() == nil {
noteFailure(&st.FreshFailures, &st.FreshFailingSince, &w.freshProgressed)
}
return w.persist(ctx, err)
}
if err := w.evaluateChanged(ctx, changes); err != nil {
if ctx.Err() != nil {
return ctx.Err()
}
if !failure(&st.FreshFailures, &st.FreshFailingSince, &w.freshProgressed) {
return w.persist(ctx, err)
}
log.Warn().Err(err).Str("org_id", w.orgID.String()).Msg("inbox tagging: changed threads failed repeatedly; page left to the cycle")
w.p.Skipped++
}
if len(changes) < w.opts.PageSize {
mark = repository.FollowUpMark{At: until, RowID: maxRowID}
} else {
last := changes[len(changes)-1]
mark = repository.FollowUpMark{At: last.At, RowID: last.RowID}
}
st.Fresh = &mark
st.FreshFailures, st.FreshFailingSince = 0, nil
w.freshProgressed = true
if err := w.save(ctx); err != nil {
return err
}
if len(changes) < w.opts.PageSize {
return nil
}
}
}
func (w *followUpWalk) evaluateChanged(ctx context.Context, changes []repository.FollowUpChange) error {
var threads []string
picked := map[string]bool{}
for _, c := range changes {
if !w.seen[c.ThreadID] && !picked[c.ThreadID] {
picked[c.ThreadID] = true
threads = append(threads, c.ThreadID)
}
}
states, err := w.s.repo.FollowUpThreadStates(ctx, w.orgID, threads, w.opts.Since)
if err != nil {
return err
}
for _, st := range states {
if err := w.evaluate(ctx, st); err != nil {
return err
}
}
return nil
}
// walk evaluates threads from `from` to the end of the last mailbox and reports whether it got there.
func (w *followUpWalk) walk(ctx context.Context, mailboxes []uuid.UUID, from *repository.FollowUpPosition) (bool, error) {
start := 0
var after *repository.FollowUpPosition
if from != nil {
for start < len(mailboxes) && bytes.Compare(mailboxes[start][:], from.MailboxID[:]) < 0 {
start++
}
if start < len(mailboxes) && mailboxes[start] == from.MailboxID {
after = from
}
}
first := true
for _, mailbox := range mailboxes[start:] {
for {
// Pages run until the first one that reads anything, so every pass moves the cycle on.
if !first && w.spent() {
return false, nil
}
if after != nil && after.At.Before(w.opts.Since) {
break
}
page, err := w.s.repo.FollowUpPage(ctx, w.orgID, mailbox, w.opts.Since, after, w.opts.PageSize)
if err != nil {
if page, err = w.pageFailed(ctx, mailbox, after, err); err != nil {
return false, err
}
}
w.p.Pages++
if page.Rows > 0 {
first = false
}
var last *repository.FollowUpPosition
for i := range page.States {
st := page.States[i]
if err := ctx.Err(); err != nil {
return false, w.stopAt(ctx, last, err)
}
if err := w.evaluate(ctx, st); err != nil {
if ctx.Err() != nil {
return false, w.stopAt(ctx, last, ctx.Err())
}
if !w.failed() {
return false, w.stopAt(ctx, last, err)
}
log.Warn().Err(err).Str("thread_id", st.ThreadID).Msg("inbox tagging: follow-up label failed repeatedly; thread skipped this cycle")
w.p.Skipped++
}
w.progressed = true
last = &st.Position
}
if page.Last != nil {
after = page.Last
if w.state != nil {
w.state.Cursor = after
w.state.PageFailures, w.state.PageFailingSince = 0, nil
w.progressed = true
if err := w.save(ctx); err != nil {
return false, err
}
}
}
if page.Rows < w.opts.PageSize {
break
}
}
after = nil
}
return true, nil
}
func (w *followUpWalk) spent() bool {
return w.opts.Budget > 0 && time.Since(w.started) >= w.opts.Budget
}
// failed counts a failure at the cycle's current place and reports whether it is time to step past it.
func (w *followUpWalk) failed() bool {
if w.state == nil {
return false
}
return failure(&w.state.PageFailures, &w.state.PageFailingSince, &w.progressed)
}
// failure counts one failure at a place and reports whether it has failed often enough, for long enough, to step past.
func failure(count *int, since **time.Time, progressed *bool) bool {
noteFailure(count, since, progressed)
if *count < followUpFailureLimit || time.Since(**since) < followUpFailureSpan {
return false
}
*count, *since = 0, nil
return true
}
// noteFailure counts one failure at a place, starting a new run when the walk has moved on since the last.
func noteFailure(count *int, since **time.Time, progressed *bool) {
if *progressed {
*count, *since, *progressed = 0, nil, false
}
if *since == nil {
now := time.Now()
*since = &now
}
*count++
}
// pageFailed steps past a page that keeps failing by reading only where it ends, or past its mailbox when even that fails.
func (w *followUpWalk) pageFailed(ctx context.Context, mailbox uuid.UUID, after *repository.FollowUpPosition, cause error) (repository.FollowUpPage, error) {
if ctx.Err() != nil {
return repository.FollowUpPage{}, cause
}
if !w.failed() {
return repository.FollowUpPage{}, w.stopAt(ctx, nil, cause)
}
page, err := w.s.repo.FollowUpPagePositions(ctx, w.orgID, mailbox, w.opts.Since, after, w.opts.PageSize)
if err != nil {
if ctx.Err() != nil {
return repository.FollowUpPage{}, err
}
log.Warn().Err(cause).Str("org_id", w.orgID.String()).Str("mailbox_id", mailbox.String()).Msg("inbox tagging: follow-up mailbox failed repeatedly; rest of it skipped this cycle")
w.p.Skipped++
return repository.FollowUpPage{Last: &repository.FollowUpPosition{MailboxID: mailbox, At: mailboxDone}}, nil
}
log.Warn().Err(cause).Str("org_id", w.orgID.String()).Str("mailbox_id", mailbox.String()).Msg("inbox tagging: follow-up page failed repeatedly; skipped this cycle")
w.p.Skipped++
return page, nil
}
// stopAt keeps the cursor on the last thread actually evaluated and returns the cause.
func (w *followUpWalk) stopAt(ctx context.Context, last *repository.FollowUpPosition, cause error) error {
if w.state == nil {
return cause
}
if last != nil {
w.state.Cursor = last
}
if w.progressed {
w.state.PageFailures, w.state.PageFailingSince = 0, nil
}
return w.persist(ctx, cause)
}
// persist saves the state on the way out of a failed pass, keeping the cause as the error.
func (w *followUpWalk) persist(ctx context.Context, cause error) error {
if err := w.saveDetached(ctx); err != nil && !errors.Is(err, errLeaseLost) {
return errors.Join(cause, err)
}
return cause
}
func (w *followUpWalk) save(ctx context.Context) error {
if w.state == nil {
return nil
}
ok, err := w.s.repo.SaveFollowUpSweep(ctx, w.orgID, w.owner, w.lease, *w.state)
if err != nil {
return err
}
if !ok {
w.lost = true
return errLeaseLost
}
return nil
}
// saveDetached saves even when the pass's context is already cancelled.
func (w *followUpWalk) saveDetached(ctx context.Context) error {
c, cancel := context.WithTimeout(context.WithoutCancel(ctx), 5*time.Second)
defer cancel()
return w.save(c)
}
func (w *followUpWalk) release(ctx context.Context) {
c, cancel := context.WithTimeout(context.WithoutCancel(ctx), 5*time.Second)
defer cancel()
if err := w.s.repo.ReleaseFollowUpSweep(c, w.orgID, w.owner); err != nil {
log.Warn().Err(err).Str("org_id", w.orgID.String()).Msg("inbox tagging: follow-up sweep lease not released")
}
}
func (w *followUpWalk) evaluate(ctx context.Context, st repository.ThreadFollowUpState) error {
if w.seen[st.ThreadID] {
return nil
}
want := FollowUp(ThreadState{
ThreadID: st.ThreadID,
LastInboundAt: st.LastInboundAt,
LastOutboundAt: st.LastOutboundAt,
BestIntent: st.BestIntent,
LastKind: st.LastKind,
}, w.now)
if err := w.s.categories.SyncExclusiveLabels(ctx, w.orgID, st.ThreadID, FollowUpLabels, want); err != nil {
return err
}
w.seen[st.ThreadID] = true
w.p.Threads++
if want == "" {
w.p.Cleared++
} else {
w.p.Labelled[want]++
}
return nil
}
+419
View File
@@ -0,0 +1,419 @@
package inboxtag
import (
"bytes"
"context"
"errors"
"fmt"
"os"
"testing"
"time"
"github.com/google/uuid"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
"github.com/warmbly/warmbly/internal/repository"
)
// The follow-up sweep against Postgres, with more threads than one old-style
// pass ever reached. Skipped unless WARMBLY_TEST_DB is set:
//
// WARMBLY_TEST_DB=postgres://warmbly:warmbly@localhost:15432/<db>?sslmode=disable \
// go test ./internal/app/inboxtag/ -run LiveFollowUpSweep -v
const liveBusyThreads = 2100
type sweepFixture struct {
t *testing.T
pool *pgxpool.Pool
owner uuid.UUID
org uuid.UUID
mailbox [2]uuid.UUID
store *repository.TagCategoryStore
}
func newSweepFixture(t *testing.T) *sweepFixture {
t.Helper()
dsn := os.Getenv("WARMBLY_TEST_DB")
if dsn == "" {
t.Skip("WARMBLY_TEST_DB not set")
}
ctx := context.Background()
pool, err := pgxpool.New(ctx, dsn)
if err != nil {
t.Fatalf("connect: %v", err)
}
t.Cleanup(pool.Close)
var version int64
if err := pool.QueryRow(ctx, `SELECT version FROM schema_migrations LIMIT 1`).Scan(&version); err != nil || version < 249 {
t.Fatalf("WARMBLY_TEST_DB is at schema version %d (err %v); this test needs 249 or later", version, err)
}
f := &sweepFixture{t: t, pool: pool, owner: uuid.New(), org: uuid.New(), mailbox: [2]uuid.UUID{uuid.New(), uuid.New()}}
f.store = repository.NewTagCategoryStore(pool)
f.exec(`INSERT INTO users (id, first_name, last_name, email, password_hash) VALUES ($1, 'Sweep', 'Owner', $2, 'x')`,
f.owner, "sweep-"+f.owner.String()[:8]+"@test.local")
f.exec(`INSERT INTO organizations (id, name, slug, owner_user_id) VALUES ($1, 'Sweep', $2, $3)`, f.org, "sweep-"+f.org.String()[:8], f.owner)
for _, mb := range f.mailbox {
f.exec(`INSERT INTO email_accounts (id, user_id, organization_id, email, name,
signature_plain, signature_html, provider, status, campaign_limit, min_wait_time, timezone)
VALUES ($1, $2, $3, $4, 'Sweep', '', '', 'smtp_imap', 'active', 50, 0, 'UTC')`,
mb, f.owner, f.org, "sweep-mb-"+mb.String()[:8]+"@test.local")
}
t.Cleanup(func() {
c := context.Background()
for _, sql := range []string{
`DELETE FROM unibox_thread_labels WHERE organization_id = $1`,
`DELETE FROM categories WHERE organization_id = $1`,
`DELETE FROM inbox_tag_results WHERE organization_id = $1`,
`DELETE FROM unibox_emails WHERE email_id IN (SELECT id FROM email_accounts WHERE organization_id = $1)`,
`DELETE FROM inbox_follow_up_sweeps WHERE organization_id = $1`,
`DELETE FROM email_accounts WHERE organization_id = $1`,
`DELETE FROM organizations WHERE id = $1`,
} {
if _, err := pool.Exec(c, sql, f.org); err != nil {
t.Errorf("cleanup %q: %v", sql, err)
}
}
if _, err := pool.Exec(c, `DELETE FROM users WHERE id = $1`, f.owner); err != nil {
t.Errorf("cleanup user: %v", err)
}
})
// The newest 2,100 threads: we wrote an hour ago, nothing is owed yet.
f.exec(`INSERT INTO unibox_emails (id, user_id, email_id, folder, provider_folder, message_id, thread_id,
from_addr, subject, body_text, internal_date)
SELECT gen_random_uuid(), $1, CASE WHEN g % 2 = 0 THEN $2::uuid ELSE $3::uuid END, 'sent', 'sent',
'<busy-' || g || '@sweep.test>', 'busy-' || g, ARRAY['me@sweep.test'], 'Hello', 'body',
NOW() - interval '1 hour' - g * interval '1 second'
FROM generate_series(1, $4::int) g`, f.owner, f.mailbox[0], f.mailbox[1], liveBusyThreads)
day := func(n int) time.Time { return time.Now().AddDate(0, 0, -n) }
// Every conversation below is older than all of those.
f.message(0, "sent", "stale-chase", day(8), "", "")
f.message(0, "sent", "stale-cold", day(12), "", "")
f.message(0, "inbox", "stale-cold", day(14), KindHumanReply, IntentAgreed)
f.message(1, "sent", "stale-owed", day(9), "", "")
f.message(1, "inbox", "stale-owed", day(4), KindHumanReply, IntentWantsInfo)
f.message(0, "sent", "stale-declined", day(20), "", "")
f.message(0, "inbox", "stale-declined", day(19), KindHumanReply, IntentNotInterested)
f.message(1, "sent", "stale-bounce", day(20), "", "")
f.message(1, "inbox", "stale-bounce", day(19), KindBounceHard, "")
f.message(0, "sent", "stale-answered", day(10), "", "")
f.message(0, "inbox", "stale-answered", day(1), KindHumanReply, IntentWantsInfo)
f.label("stale-answered", LabelFollowUp)
// One conversation held by both mailboxes is still one thread.
f.message(1, "inbox", "stale-shared", day(9), "", "")
f.message(0, "sent", "stale-shared", day(8), "", "")
return f
}
// liveThreads is every thread the fixture has that we wrote to.
const liveThreads = liveBusyThreads + 7
func (f *sweepFixture) exec(sql string, args ...any) {
f.t.Helper()
if _, err := f.pool.Exec(context.Background(), sql, args...); err != nil {
f.t.Fatalf("fixture %q: %v", sql[:min(70, len(sql))], err)
}
}
func (f *sweepFixture) message(mailbox int, folder, thread string, at time.Time, kind, intent string) {
f.t.Helper()
id := fmt.Sprintf("<%s-%s-%d@sweep.test>", thread, folder, mailbox)
f.exec(`INSERT INTO unibox_emails (id, user_id, email_id, folder, provider_folder, message_id, thread_id,
from_addr, subject, body_text, internal_date)
VALUES ($1, $2, $3, $4, $4, $5, $6, ARRAY['x@sweep.test'], 'Hello', 'body', $7)`,
uuid.New(), f.owner, f.mailbox[mailbox], folder, id, thread, at)
if kind != "" {
f.exec(`INSERT INTO inbox_tag_results (organization_id, email_account_id, message_id, thread_id, status, kind, intent)
VALUES ($1, $2, $3, $4, 'complete', $5, $6)`, f.org, f.mailbox[mailbox], id, thread, kind, intent)
}
}
func (f *sweepFixture) label(thread, slug string) {
f.t.Helper()
id, err := f.store.EnsureCategory(context.Background(), f.org, slug)
if err != nil {
f.t.Fatalf("category: %v", err)
}
f.exec(`INSERT INTO unibox_thread_labels (organization_id, thread_id, category_id) VALUES ($1, $2, $3)`, f.org, thread, id)
}
func (f *sweepFixture) labels(thread string) []string {
f.t.Helper()
rows, err := f.pool.Query(context.Background(), `
SELECT c.title FROM unibox_thread_labels l JOIN categories c ON c.id = l.category_id
WHERE l.organization_id = $1 AND l.thread_id = $2 ORDER BY c.title`, f.org, thread)
if err != nil {
f.t.Fatalf("labels: %v", err)
}
defer rows.Close()
var out []string
for rows.Next() {
var title string
if err := rows.Scan(&title); err != nil {
f.t.Fatalf("scan: %v", err)
}
out = append(out, title)
}
return out
}
// pass is one hourly pass from a freshly started service, as after a consumer restart.
func (f *sweepFixture) pass() FollowUpProgress {
f.t.Helper()
return f.passWith(FollowUpSweep{Since: time.Now().AddDate(0, 0, -90), Budget: time.Nanosecond, PageSize: 500})
}
func (f *sweepFixture) passWith(opts FollowUpSweep) FollowUpProgress {
f.t.Helper()
repo := repository.NewInboxTagRepository(f.pool)
svc := NewService(nil, repo, repository.NewTagCategoryStore(f.pool), nil, true)
p, err := svc.SweepFollowUps(context.Background(), f.org, opts)
if err != nil {
f.t.Fatalf("sweep: %v", err)
}
return p
}
func (f *sweepFixture) cursor() *repository.FollowUpPosition {
f.t.Helper()
var mailbox, row *uuid.UUID
var at *time.Time
err := f.pool.QueryRow(context.Background(), `
SELECT email_account_id, internal_date, message_row_id FROM inbox_follow_up_sweeps WHERE organization_id = $1`,
f.org).Scan(&mailbox, &at, &row)
if errors.Is(err, pgx.ErrNoRows) {
return nil
}
if err != nil {
f.t.Fatalf("cursor: %v", err)
}
if mailbox == nil || at == nil || row == nil {
return nil
}
return &repository.FollowUpPosition{MailboxID: *mailbox, At: *at, RowID: *row}
}
// A conversation older than the newest 2,000 is reached within one cycle and
// gets the label its calendar says, with automated and declined threads left
// alone and a stale label taken off.
func TestLiveFollowUpSweepReachesThreadsBeyondTheNewest2000(t *testing.T) {
f := newSweepFixture(t)
threads := 0
for passes := 1; ; passes++ {
p := f.pass()
threads += p.Threads
if p.Complete {
break
}
if passes > 20 {
t.Fatal("the cycle never completed")
}
}
if threads != liveThreads {
t.Errorf("a cycle evaluated %d threads, want each of the %d once", threads, liveThreads)
}
for thread, want := range map[string]string{
"stale-chase": LabelFollowUp,
"stale-cold": LabelGoneQuiet,
"stale-owed": LabelNeedsReply,
"stale-declined": "",
"stale-bounce": "",
"stale-answered": "",
"stale-shared": LabelFollowUp,
"busy-1": "",
} {
got := f.labels(thread)
switch {
case want == "" && len(got) != 0:
t.Errorf("%s wears %v, want nothing", thread, got)
case want != "" && (len(got) != 1 || got[0] != want):
t.Errorf("%s wears %v, want %q", thread, got, want)
}
}
if pos := f.cursor(); pos != nil {
t.Errorf("a finished cycle left a cursor at %+v", pos)
}
}
// A pass that stops resumes where it left off from the saved cursor, so each
// thread is evaluated once per cycle rather than the newest page every time.
func TestLiveFollowUpSweepResumesFromItsCursor(t *testing.T) {
f := newSweepFixture(t)
first := f.pass()
if first.Complete || first.Threads == 0 {
t.Fatalf("first pass %+v, want one partial page", first)
}
prev := f.cursor()
if prev == nil {
t.Fatal("a pass stopped mid-cycle without saving where")
}
threads := first.Threads
for passes := 2; ; passes++ {
p := f.pass()
threads += p.Threads
if p.Complete {
break
}
pos := f.cursor()
if pos == nil || !walkedPast(*pos, *prev) {
t.Fatalf("pass %d left the cursor at %+v, not past %+v", passes, pos, prev)
}
prev = pos
if passes > 20 {
t.Fatal("the cycle never completed")
}
}
if threads != liveThreads {
t.Errorf("the resumed passes evaluated %d threads, want each of the %d once", threads, liveThreads)
}
if got := f.labels("stale-chase"); len(got) != 1 || got[0] != LabelFollowUp {
t.Errorf("stale-chase wears %v after a resumed cycle", got)
}
}
// walkedPast reports whether a comes after b in the sweep's walk order.
func walkedPast(a, b repository.FollowUpPosition) bool {
if c := bytes.Compare(a.MailboxID[:], b.MailboxID[:]); c != 0 {
return c > 0
}
if !a.At.Equal(b.At) {
return a.At.Before(b.At)
}
return bytes.Compare(a.RowID[:], b.RowID[:]) < 0
}
// A reply the sync stored late, under a date older than any fixed window, and
// a verdict written long after its message, are both checked in the next pass
// while the cycle is still among the newest threads.
func TestLiveFollowUpSweepChecksLateSyncedReplies(t *testing.T) {
f := newSweepFixture(t)
ctx := context.Background()
// Everything the fixture wrote was stored and judged two days ago.
f.exec(`UPDATE unibox_emails SET ingested_at = NOW() - interval '2 days'
WHERE email_id IN (SELECT id FROM email_accounts WHERE organization_id = $1)`, f.org)
f.exec(`UPDATE inbox_tag_results SET updated_at = NOW() - interval '2 days' WHERE organization_id = $1`, f.org)
f.label("stale-chase", LabelFollowUp)
// Inbound four days ago, after our send, not judged yet: nothing is owed until it is.
f.message(1, "sent", "late-verdict", time.Now().AddDate(0, 0, -9), "", "")
f.message(1, "inbox", "late-verdict", time.Now().AddDate(0, 0, -4), "", "")
f.exec(`UPDATE unibox_emails SET ingested_at = NOW() - interval '2 days' WHERE thread_id = 'late-verdict'`)
opts := FollowUpSweep{Since: time.Now().AddDate(0, 0, -90), Fresh: 24 * time.Hour, Budget: time.Nanosecond, PageSize: 500}
if p := f.passWith(opts); p.Complete {
t.Fatalf("first pass finished the cycle; the test needs it elsewhere: %+v", p)
}
if got := f.labels("late-verdict"); len(got) != 0 {
t.Fatalf("late-verdict wears %v before its reply was judged", got)
}
// As if that pass ran ten minutes ago.
f.exec(`UPDATE inbox_follow_up_sweeps SET fresh_at = fresh_at - interval '10 minutes' WHERE organization_id = $1`, f.org)
// They answered six hours ago; the sync stored it five minutes ago.
id := "<late-reply@sweep.test>"
f.exec(`INSERT INTO unibox_emails (id, user_id, email_id, folder, provider_folder, message_id, thread_id,
from_addr, subject, body_text, internal_date, ingested_at)
VALUES ($1, $2, $3, 'inbox', 'inbox', $4, 'stale-chase', ARRAY['x@sweep.test'], 'Re: Hello', 'body',
NOW() - interval '6 hours', NOW() - interval '5 minutes')`, uuid.New(), f.owner, f.mailbox[0], id)
f.exec(`INSERT INTO inbox_tag_results (organization_id, email_account_id, message_id, thread_id, status, kind, intent, updated_at)
VALUES ($1, $2, $3, 'stale-chase', 'complete', $4, $5, NOW() - interval '2 days')`,
f.org, f.mailbox[0], id, KindHumanReply, IntentWantsInfo)
// The four-day-old reply is judged five minutes ago.
f.exec(`INSERT INTO inbox_tag_results (organization_id, email_account_id, message_id, thread_id, status, kind, intent, updated_at)
VALUES ($1, $2, '<late-verdict-inbox-1@sweep.test>', 'late-verdict', 'complete', $3, $4, NOW() - interval '5 minutes')`,
f.org, f.mailbox[1], KindHumanReply, IntentWantsInfo)
before := f.cursor()
f.passWith(opts)
if got := f.labels("stale-chase"); len(got) != 0 {
t.Errorf("stale-chase still wears %v after a reply stored since the last pass", got)
}
if got := f.labels("late-verdict"); len(got) != 1 || got[0] != LabelNeedsReply {
t.Errorf("late-verdict wears %v after its reply was judged, want Needs reply", got)
}
if after := f.cursor(); after == nil || before == nil || after.MailboxID != before.MailboxID {
t.Fatalf("the cycle left the first mailbox (%+v -> %+v); the threads may have been reached by it", before, after)
}
var freshAt time.Time
if err := f.pool.QueryRow(ctx, `SELECT fresh_at FROM inbox_follow_up_sweeps WHERE organization_id = $1`, f.org).Scan(&freshAt); err != nil {
t.Fatalf("fresh_at: %v", err)
}
if time.Since(freshAt) > 3*time.Minute {
t.Errorf("the changed-thread mark stayed at %v", freshAt)
}
}
// While one walker holds a workspace, another neither sweeps it nor moves its state.
func TestLiveFollowUpSweepLeaseKeepsOneWalker(t *testing.T) {
f := newSweepFixture(t)
ctx := context.Background()
f.pass()
before := f.cursor()
if before == nil {
t.Fatal("first pass saved no cursor")
}
repo := repository.NewInboxTagRepository(f.pool)
holder := uuid.New()
st, err := repo.ClaimFollowUpSweep(ctx, f.org, holder, time.Minute)
if err != nil || st == nil {
t.Fatalf("claim: %+v %v", st, err)
}
if p := f.pass(); !p.Busy || p.Threads != 0 {
t.Fatalf("a second walker swept a leased workspace: %+v", p)
}
if again, err := repo.ClaimFollowUpSweep(ctx, f.org, uuid.New(), time.Minute); err != nil || again != nil {
t.Fatalf("a second claim succeeded: %+v %v", again, err)
}
if ok, err := repo.SaveFollowUpSweep(ctx, f.org, uuid.New(), time.Minute, repository.FollowUpSweepState{}); err != nil || ok {
t.Fatalf("a walker without the lease saved: %v %v", ok, err)
}
if got := f.cursor(); got == nil || !got.At.Equal(before.At) || got.RowID != before.RowID {
t.Fatalf("the cursor moved to %+v while another walker held the lease", got)
}
if ok, err := repo.SaveFollowUpSweep(ctx, f.org, holder, time.Minute, *st); err != nil || !ok {
t.Fatalf("the holder could not save: %v %v", ok, err)
}
if err := repo.ReleaseFollowUpSweep(ctx, f.org, holder); err != nil {
t.Fatalf("release: %v", err)
}
if p := f.pass(); p.Busy || p.Threads == 0 {
t.Fatalf("the released workspace was not swept: %+v", p)
}
}
// A walker that outlives its lease keeps saving until another walker takes the workspace over.
func TestLiveFollowUpSweepLeaseOverrun(t *testing.T) {
f := newSweepFixture(t)
ctx := context.Background()
repo := repository.NewInboxTagRepository(f.pool)
slow := uuid.New()
st, err := repo.ClaimFollowUpSweep(ctx, f.org, slow, time.Millisecond)
if err != nil || st == nil {
t.Fatalf("claim: %+v %v", st, err)
}
time.Sleep(20 * time.Millisecond)
if ok, err := repo.SaveFollowUpSweep(ctx, f.org, slow, time.Millisecond, *st); err != nil || !ok {
t.Fatalf("the holder could not save after its lease lapsed: %v %v", ok, err)
}
time.Sleep(20 * time.Millisecond)
next := uuid.New()
taken, err := repo.ClaimFollowUpSweep(ctx, f.org, next, time.Minute)
if err != nil || taken == nil {
t.Fatalf("a lapsed lease could not be taken over: %+v %v", taken, err)
}
if ok, err := repo.SaveFollowUpSweep(ctx, f.org, slow, time.Minute, *st); err != nil || ok {
t.Fatalf("the old holder saved after a takeover: %v %v", ok, err)
}
if ok, err := repo.SaveFollowUpSweep(ctx, f.org, next, time.Minute, *taken); err != nil || !ok {
t.Fatalf("the new holder could not save: %v %v", ok, err)
}
}
+1
View File
@@ -926,6 +926,7 @@ var ExcludedTables = map[string]string{
"contact_imports": "Contact imports in progress or recently finished. They are work this instance is doing; the contacts they created travel with the contacts group.",
"contact_import_rows": "The uploaded rows of a contact import and what became of each. They follow contact_imports, which does not travel.",
"placement_renders": "The copy a tracking comparison is sending to each seed, sealed so both halves send the same words. It lives only while the comparison runs, and a copy that had not been sent stays behind with its task.",
"inbox_follow_up_sweeps": "This instance's hourly follow-up sweep state for the workspace: where its cycle stopped (by this instance's mailbox and message row ids), how far it has checked changed conversations, and which walker holds it. The destination starts its own cycle at the newest conversation.",
"user_view_preferences": "Each member's own column layout and sort for the dashboard's lists, and their unibox scope rail arrangement. It belongs to the person rather than the workspace: members are matched by account on import and a layout names custom fields the destination may not hold yet, so everyone starts from the default view and picks their columns again.",
}
@@ -0,0 +1 @@
DROP TABLE IF EXISTS public.inbox_follow_up_sweeps;
@@ -0,0 +1,19 @@
-- Each workspace's follow-up sweep bookkeeping: cycle cursor, changed-thread watermark, failure counts and walker lease.
CREATE TABLE IF NOT EXISTS public.inbox_follow_up_sweeps (
organization_id uuid PRIMARY KEY REFERENCES public.organizations(id) ON DELETE CASCADE,
-- No foreign key: a mailbox deleted mid-cycle still orders the walk. NULL starts a cycle at the newest.
email_account_id uuid,
internal_date timestamptz,
message_row_id uuid,
-- Changes up to (fresh_at, fresh_row_id) have been checked.
fresh_at timestamptz,
fresh_row_id uuid,
page_failures integer NOT NULL DEFAULT 0,
fresh_failures integer NOT NULL DEFAULT 0,
-- When the current run of failures began; a place is stepped past only after failing for a while.
page_failing_since timestamptz,
fresh_failing_since timestamptz,
lease_owner uuid,
leased_until timestamptz,
updated_at timestamptz NOT NULL DEFAULT NOW()
);
@@ -0,0 +1 @@
DROP INDEX CONCURRENTLY IF EXISTS idx_unibox_emails_account_date;
@@ -0,0 +1,3 @@
-- Alone in its file for CONCURRENTLY. The follow-up sweep pages one mailbox newest first on (internal_date, id).
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_unibox_emails_account_date
ON public.unibox_emails USING btree (email_id, internal_date, id);
@@ -0,0 +1 @@
ALTER TABLE public.unibox_emails DROP COLUMN IF EXISTS ingested_at;
@@ -0,0 +1,3 @@
-- When the row was stored (an imported row keeps its source value); created_at is the worker's sync clock. Existing rows read as epoch.
ALTER TABLE public.unibox_emails ADD COLUMN IF NOT EXISTS ingested_at timestamptz NOT NULL DEFAULT 'epoch';
ALTER TABLE public.unibox_emails ALTER COLUMN ingested_at SET DEFAULT NOW();
@@ -0,0 +1 @@
DROP INDEX CONCURRENTLY IF EXISTS idx_unibox_emails_account_ingested;
@@ -0,0 +1,3 @@
-- Alone in its file for CONCURRENTLY. The follow-up sweep reads each mailbox's newly stored messages in order.
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_unibox_emails_account_ingested
ON public.unibox_emails USING btree (email_id, ingested_at, id);
@@ -0,0 +1 @@
DROP INDEX CONCURRENTLY IF EXISTS idx_inbox_tag_results_updated;
@@ -0,0 +1,3 @@
-- Alone in its file for CONCURRENTLY. The follow-up sweep reads a workspace's newly written verdicts in order.
CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_inbox_tag_results_updated
ON public.inbox_tag_results USING btree (organization_id, updated_at, id);
+384
View File
@@ -0,0 +1,384 @@
package repository
import (
"context"
"errors"
"time"
"github.com/google/uuid"
"github.com/jackc/pgx/v5"
)
// ThreadFollowUpState is one thread's follow-up facts. Every field is read from
// the database; none of it is inferred, and none of it is asked of a model.
type ThreadFollowUpState struct {
ThreadID string
LastInboundAt time.Time
LastOutboundAt time.Time
// BestIntent is the most recent trusted intent in this thread.
BestIntent string
// LastKind is the classified kind of the newest inbound message, which is
// what says whether the "reply" was a person or a mail server.
LastKind string
// Position is the thread's newest message on a cycle page; zero on the changed-thread check.
Position FollowUpPosition
}
// FollowUpPosition is one message row in the walk: mailboxes by id, then (internal_date, id) descending.
type FollowUpPosition struct {
MailboxID uuid.UUID
At time.Time
RowID uuid.UUID
}
// FollowUpPage is one bounded step of the walk.
type FollowUpPage struct {
// States are the page's threads that we have written to, in walk order.
States []ThreadFollowUpState
// Rows is how many messages the page read; fewer than the limit ends the mailbox.
Rows int
// Last is the last message read, nil when the page was empty.
Last *FollowUpPosition
}
// FollowUpMark is a place in the stream of stored messages and written verdicts, ordered by (At, RowID).
type FollowUpMark struct {
At time.Time
RowID uuid.UUID
}
// FollowUpChange is one message stored or verdict written.
type FollowUpChange struct {
At time.Time
RowID uuid.UUID
ThreadID string
}
// FollowUpSweepState is a workspace's sweep bookkeeping as its walker claimed it.
type FollowUpSweepState struct {
// Cursor is where the cycle stopped; nil starts a cycle at the newest.
Cursor *FollowUpPosition
// Fresh is how far changed threads have been checked; nil when never.
Fresh *FollowUpMark
PageFailures int
FreshFailures int
// PageFailingSince and FreshFailingSince are when the current run of failures began.
PageFailingSince *time.Time
FreshFailingSince *time.Time
// Now is the database clock at the claim.
Now time.Time
}
func (r *inboxTagRepository) FollowUpMailboxes(ctx context.Context, orgID uuid.UUID) ([]uuid.UUID, error) {
rows, err := r.db.Query(ctx, `SELECT id FROM email_accounts WHERE organization_id = $1 ORDER BY id`, orgID)
if err != nil {
return nil, err
}
defer rows.Close()
var out []uuid.UUID
for rows.Next() {
var id uuid.UUID
if err := rows.Scan(&id); err != nil {
return nil, err
}
out = append(out, id)
}
return out, rows.Err()
}
// followUpPageSQL is one mailbox's next messages, newest first after an optional keyset; $1 org, $2 mailbox, $3 since, $4 limit.
func followUpPageSQL(keyset string) string {
return `page AS MATERIALIZED (
SELECT ue.id, ue.internal_date, ue.thread_id
FROM unibox_emails ue
WHERE ue.email_id = $2
AND EXISTS (SELECT 1 FROM email_accounts ea WHERE ea.id = $2 AND ea.organization_id = $1)
AND ue.internal_date >= $3 ` + keyset + `
AND ue.thread_id <> ''
ORDER BY ue.internal_date DESC, ue.id DESC
LIMIT $4
)`
}
// followUpStatesSQL resolves follow-up facts for the threads in a heads(id, thread_id) CTE; $1 is the org.
const followUpStatesSQL = `states AS (
SELECT h.id, h.thread_id, agg.last_in, agg.last_out, agg.last_any,
COALESCE(best.intent, '') AS intent, COALESCE(newest.kind, '') AS kind
FROM heads h
CROSS JOIN LATERAL (
SELECT MAX(o.internal_date) FILTER (WHERE o.folder = 'inbox') AS last_in,
MAX(o.internal_date) FILTER (WHERE o.folder = 'sent') AS last_out,
MAX(o.internal_date) AS last_any
FROM unibox_emails o
JOIN email_accounts oa ON oa.id = o.email_id AND oa.organization_id = $1
WHERE o.thread_id = h.thread_id
) agg
LEFT JOIN LATERAL (
SELECT o.email_id, o.message_id
FROM unibox_emails o
JOIN email_accounts oa ON oa.id = o.email_id AND oa.organization_id = $1
WHERE o.thread_id = h.thread_id AND o.folder = 'inbox'
ORDER BY o.internal_date DESC, o.id DESC
LIMIT 1
) latest ON TRUE
LEFT JOIN LATERAL (
SELECT r.intent
FROM inbox_tag_results r
JOIN unibox_emails ue
ON ue.email_id = r.email_account_id
AND ue.thread_id = r.thread_id
AND ue.message_id = r.message_id
JOIN email_accounts oa ON oa.id = ue.email_id AND oa.organization_id = $1
WHERE r.organization_id = $1 AND r.thread_id = h.thread_id
AND r.status = 'complete' AND r.review_reason <> 'intent' AND r.intent <> ''
ORDER BY ue.internal_date DESC, r.created_at DESC
LIMIT 1
) best ON TRUE
LEFT JOIN inbox_tag_results newest
ON newest.organization_id = $1
AND newest.email_account_id = latest.email_id
AND newest.thread_id = h.thread_id
AND newest.message_id = latest.message_id
AND newest.status = 'complete'
AND newest.review_reason <> 'kind'
WHERE agg.last_out IS NOT NULL
)`
// FollowUpPage reads one bounded page of a mailbox and resolves facts only for threads whose newest workspace message is on it.
func (r *inboxTagRepository) FollowUpPage(ctx context.Context, orgID, mailboxID uuid.UUID, since time.Time, after *FollowUpPosition, limit int) (FollowUpPage, error) {
args, keyset := followUpPageArgs(orgID, mailboxID, since, after, limit)
q := `WITH ` + followUpPageSQL(keyset) + `,
heads AS (
SELECT p.id, p.thread_id
FROM page p
WHERE NOT EXISTS (
SELECT 1
FROM unibox_emails o
JOIN email_accounts oa ON oa.id = o.email_id AND oa.organization_id = $1
WHERE o.thread_id = p.thread_id
AND (o.internal_date, o.id) > (p.internal_date, p.id)
)
),
` + followUpStatesSQL + `
SELECT p.id, p.internal_date, p.thread_id, s.id IS NOT NULL,
s.last_in, s.last_out, COALESCE(s.intent, ''), COALESCE(s.kind, '')
FROM page p
LEFT JOIN states s ON s.id = p.id
ORDER BY p.internal_date DESC, p.id DESC
`
rows, err := r.db.Query(ctx, q, args...)
if err != nil {
return FollowUpPage{}, err
}
defer rows.Close()
var out FollowUpPage
for rows.Next() {
var (
pos = FollowUpPosition{MailboxID: mailboxID}
threadID string
due bool
lastIn, lastOut *time.Time
intent, kind string
)
if err := rows.Scan(&pos.RowID, &pos.At, &threadID, &due, &lastIn, &lastOut, &intent, &kind); err != nil {
return FollowUpPage{}, err
}
out.Rows++
last := pos
out.Last = &last
if due {
out.States = append(out.States, followUpState(threadID, lastIn, lastOut, intent, kind, pos))
}
}
return out, rows.Err()
}
// FollowUpPagePositions reads only where a page ends, to step past a page whose facts cannot be read.
func (r *inboxTagRepository) FollowUpPagePositions(ctx context.Context, orgID, mailboxID uuid.UUID, since time.Time, after *FollowUpPosition, limit int) (FollowUpPage, error) {
args, keyset := followUpPageArgs(orgID, mailboxID, since, after, limit)
rows, err := r.db.Query(ctx, `WITH `+followUpPageSQL(keyset)+`
SELECT id, internal_date FROM page ORDER BY internal_date DESC, id DESC`, args...)
if err != nil {
return FollowUpPage{}, err
}
defer rows.Close()
var out FollowUpPage
for rows.Next() {
pos := FollowUpPosition{MailboxID: mailboxID}
if err := rows.Scan(&pos.RowID, &pos.At); err != nil {
return FollowUpPage{}, err
}
out.Rows++
out.Last = &pos
}
return out, rows.Err()
}
func followUpPageArgs(orgID, mailboxID uuid.UUID, since time.Time, after *FollowUpPosition, limit int) ([]any, string) {
args := []any{orgID, mailboxID, since, limit}
if after == nil {
return args, ""
}
return append(args, after.At, after.RowID), `AND (ue.internal_date, ue.id) < ($5, $6)`
}
func followUpState(threadID string, lastIn, lastOut *time.Time, intent, kind string, pos FollowUpPosition) ThreadFollowUpState {
st := ThreadFollowUpState{ThreadID: threadID, BestIntent: intent, LastKind: kind, Position: pos}
if lastIn != nil {
st.LastInboundAt = *lastIn
}
if lastOut != nil {
st.LastOutboundAt = *lastOut
}
return st
}
// FollowUpChanges reads the next messages stored and verdicts written in a workspace after a mark, up to until.
func (r *inboxTagRepository) FollowUpChanges(ctx context.Context, orgID uuid.UUID, after FollowUpMark, until time.Time, limit int) ([]FollowUpChange, error) {
rows, err := r.db.Query(ctx, `
SELECT c.at, c.id, c.thread_id
FROM (
SELECT m.at, m.id, m.thread_id
FROM email_accounts ea
CROSS JOIN LATERAL (
SELECT ue.ingested_at AS at, ue.id, ue.thread_id
FROM unibox_emails ue
WHERE ue.email_id = ea.id
AND (ue.ingested_at, ue.id) > ($2, $3) AND ue.ingested_at <= $4
AND ue.thread_id <> ''
ORDER BY ue.ingested_at, ue.id
LIMIT $5
) m
WHERE ea.organization_id = $1
UNION ALL
(SELECT r.updated_at, r.id, r.thread_id
FROM inbox_tag_results r
WHERE r.organization_id = $1
AND (r.updated_at, r.id) > ($2, $3) AND r.updated_at <= $4
AND r.thread_id <> ''
ORDER BY r.updated_at, r.id
LIMIT $5)
) c
ORDER BY c.at, c.id
LIMIT $5`, orgID, after.At, after.RowID, until, limit)
if err != nil {
return nil, err
}
defer rows.Close()
var out []FollowUpChange
for rows.Next() {
var c FollowUpChange
if err := rows.Scan(&c.At, &c.RowID, &c.ThreadID); err != nil {
return nil, err
}
out = append(out, c)
}
return out, rows.Err()
}
// FollowUpThreadStates resolves follow-up facts for named threads with activity since a cutoff.
func (r *inboxTagRepository) FollowUpThreadStates(ctx context.Context, orgID uuid.UUID, threadIDs []string, since time.Time) ([]ThreadFollowUpState, error) {
if len(threadIDs) == 0 {
return nil, nil
}
rows, err := r.db.Query(ctx, `
WITH heads AS (
SELECT DISTINCT NULL::uuid AS id, t AS thread_id FROM unnest($2::text[]) AS t WHERE t <> ''
),
`+followUpStatesSQL+`
SELECT s.thread_id, s.last_in, s.last_out, s.intent, s.kind
FROM states s
WHERE s.last_any >= $3
ORDER BY s.thread_id`, orgID, threadIDs, since)
if err != nil {
return nil, err
}
defer rows.Close()
var out []ThreadFollowUpState
for rows.Next() {
var (
threadID, intent, kind string
lastIn, lastOut *time.Time
)
if err := rows.Scan(&threadID, &lastIn, &lastOut, &intent, &kind); err != nil {
return nil, err
}
out = append(out, followUpState(threadID, lastIn, lastOut, intent, kind, FollowUpPosition{}))
}
return out, rows.Err()
}
// ClaimFollowUpSweep takes the workspace's sweep lease and returns its state; nil while another walker holds it.
func (r *inboxTagRepository) ClaimFollowUpSweep(ctx context.Context, orgID, owner uuid.UUID, lease time.Duration) (*FollowUpSweepState, error) {
var (
st FollowUpSweepState
mailbox, row, freshID *uuid.UUID
at, freshAt *time.Time
)
err := r.db.QueryRow(ctx, `
INSERT INTO inbox_follow_up_sweeps (organization_id, lease_owner, leased_until, updated_at)
VALUES ($1, $2, NOW() + make_interval(secs => $3), NOW())
ON CONFLICT (organization_id) DO UPDATE SET
lease_owner = EXCLUDED.lease_owner,
leased_until = EXCLUDED.leased_until,
updated_at = NOW()
WHERE inbox_follow_up_sweeps.leased_until IS NULL
OR inbox_follow_up_sweeps.leased_until <= NOW()
OR inbox_follow_up_sweeps.lease_owner = EXCLUDED.lease_owner
RETURNING email_account_id, internal_date, message_row_id, fresh_at, fresh_row_id,
page_failures, fresh_failures, page_failing_since, fresh_failing_since, NOW()`, orgID, owner, lease.Seconds(),
).Scan(&mailbox, &at, &row, &freshAt, &freshID, &st.PageFailures, &st.FreshFailures,
&st.PageFailingSince, &st.FreshFailingSince, &st.Now)
if errors.Is(err, pgx.ErrNoRows) {
return nil, nil
}
if err != nil {
return nil, err
}
if mailbox != nil && at != nil && row != nil {
st.Cursor = &FollowUpPosition{MailboxID: *mailbox, At: *at, RowID: *row}
}
if freshAt != nil {
st.Fresh = &FollowUpMark{At: *freshAt}
if freshID != nil {
st.Fresh.RowID = *freshID
}
}
return &st, nil
}
// SaveFollowUpSweep stores the state and renews the lease; false once another walker has taken it over.
func (r *inboxTagRepository) SaveFollowUpSweep(ctx context.Context, orgID, owner uuid.UUID, lease time.Duration, st FollowUpSweepState) (bool, error) {
var mailbox, row, freshID *uuid.UUID
var at, freshAt *time.Time
if st.Cursor != nil {
mailbox, at, row = &st.Cursor.MailboxID, &st.Cursor.At, &st.Cursor.RowID
}
if st.Fresh != nil {
freshAt, freshID = &st.Fresh.At, &st.Fresh.RowID
}
tag, err := r.db.Exec(ctx, `
UPDATE inbox_follow_up_sweeps SET
email_account_id = $3, internal_date = $4, message_row_id = $5,
fresh_at = $6, fresh_row_id = $7,
page_failures = $8, fresh_failures = $9,
page_failing_since = $10, fresh_failing_since = $11,
leased_until = NOW() + make_interval(secs => $12),
updated_at = NOW()
WHERE organization_id = $1 AND lease_owner = $2`,
orgID, owner, mailbox, at, row, freshAt, freshID, st.PageFailures, st.FreshFailures,
st.PageFailingSince, st.FreshFailingSince, lease.Seconds())
if err != nil {
return false, err
}
return tag.RowsAffected() == 1, nil
}
// ReleaseFollowUpSweep gives the lease back so the next pass need not wait for it to lapse.
func (r *inboxTagRepository) ReleaseFollowUpSweep(ctx context.Context, orgID, owner uuid.UUID) error {
_, err := r.db.Exec(ctx, `
UPDATE inbox_follow_up_sweeps SET leased_until = NULL, updated_at = NOW()
WHERE organization_id = $1 AND lease_owner = $2`, orgID, owner)
return err
}
+10 -101
View File
@@ -67,9 +67,16 @@ type InboxTagRepository interface {
// before they were asked whether they need acting on.
ListUncheckedNotifications(ctx context.Context, orgID uuid.UUID, since time.Time, limit int) ([]BackfillCandidate, error)
// ThreadStates backs the follow-up sweep: who spoke last, when, and how far
// the thread ever got.
ThreadStates(ctx context.Context, orgID uuid.UUID, since time.Time, limit int) ([]ThreadFollowUpState, error)
// The FollowUp* methods back the follow-up sweep: a cycle that walks each
// mailbox newest first a page at a time, and a check of changed threads.
FollowUpMailboxes(ctx context.Context, orgID uuid.UUID) ([]uuid.UUID, error)
FollowUpPage(ctx context.Context, orgID, mailboxID uuid.UUID, since time.Time, after *FollowUpPosition, limit int) (FollowUpPage, error)
FollowUpPagePositions(ctx context.Context, orgID, mailboxID uuid.UUID, since time.Time, after *FollowUpPosition, limit int) (FollowUpPage, error)
FollowUpChanges(ctx context.Context, orgID uuid.UUID, after FollowUpMark, until time.Time, limit int) ([]FollowUpChange, error)
FollowUpThreadStates(ctx context.Context, orgID uuid.UUID, threadIDs []string, since time.Time) ([]ThreadFollowUpState, error)
ClaimFollowUpSweep(ctx context.Context, orgID, owner uuid.UUID, lease time.Duration) (*FollowUpSweepState, error)
SaveFollowUpSweep(ctx context.Context, orgID, owner uuid.UUID, lease time.Duration, st FollowUpSweepState) (bool, error)
ReleaseFollowUpSweep(ctx context.Context, orgID, owner uuid.UUID) error
// GetByMessageID reads one completed verdict, so the reply classifier, the
// inbox agent and the action executor can reuse a judgment already paid
@@ -485,101 +492,3 @@ func (r *inboxTagRepository) PreviousOutbound(ctx context.Context, accountID uui
}
return body, campaign, nil
}
// ThreadFollowUpState is one thread's follow-up facts. Every field is read from
// the database; none of it is inferred, and none of it is asked of a model.
type ThreadFollowUpState struct {
ThreadID string
LastInboundAt time.Time
LastOutboundAt time.Time
// BestIntent is the most recent trusted intent in this thread.
BestIntent string
// LastKind is the classified kind of the newest inbound message, which is
// what says whether the "reply" was a person or a mail server.
LastKind string
}
// ThreadStates returns follow-up facts for every thread with activity since a
// cutoff.
//
// Scoped through email_accounts because unibox_emails carries no organization
// of its own. Follow-up labels use the same organization-plus-thread key as the
// rest of the unibox.
func (r *inboxTagRepository) ThreadStates(ctx context.Context, orgID uuid.UUID, since time.Time, limit int) ([]ThreadFollowUpState, error) {
const q = `
WITH scoped_emails AS (
SELECT ue.*
FROM unibox_emails ue
JOIN email_accounts ea ON ea.id = ue.email_id
WHERE ea.organization_id = $1 AND ue.thread_id <> ''
),
active_threads AS (
SELECT DISTINCT thread_id
FROM scoped_emails
WHERE internal_date >= $2
),
threads AS (
SELECT ue.thread_id,
MAX(ue.internal_date) FILTER (WHERE ue.folder = 'inbox') AS last_in,
MAX(ue.internal_date) FILTER (WHERE ue.folder = 'sent') AS last_out
FROM scoped_emails ue
JOIN active_threads active ON active.thread_id = ue.thread_id
GROUP BY ue.thread_id
),
latest_inbound AS (
SELECT DISTINCT ON (ue.thread_id)
ue.thread_id, ue.email_id, ue.message_id
FROM scoped_emails ue
JOIN active_threads active ON active.thread_id = ue.thread_id
WHERE ue.folder = 'inbox'
ORDER BY ue.thread_id, ue.internal_date DESC
)
SELECT t.thread_id, t.last_in, t.last_out,
COALESCE(best.intent, ''), COALESCE(newest.kind, '')
FROM threads t
LEFT JOIN LATERAL (
SELECT r.intent
FROM inbox_tag_results r
JOIN scoped_emails ue
ON ue.email_id = r.email_account_id
AND ue.thread_id = r.thread_id
AND ue.message_id = r.message_id
WHERE r.organization_id = $1 AND r.thread_id = t.thread_id
AND r.status = 'complete' AND r.review_reason <> 'intent' AND r.intent <> ''
ORDER BY ue.internal_date DESC, r.created_at DESC LIMIT 1
) best ON TRUE
LEFT JOIN latest_inbound latest ON latest.thread_id = t.thread_id
LEFT JOIN inbox_tag_results newest
ON newest.organization_id = $1
AND newest.email_account_id = latest.email_id
AND newest.thread_id = latest.thread_id
AND newest.message_id = latest.message_id
AND newest.status = 'complete'
AND newest.review_reason <> 'kind'
WHERE t.last_out IS NOT NULL
ORDER BY GREATEST(COALESCE(t.last_in, 'epoch'::timestamptz), t.last_out) DESC
LIMIT $3
`
rows, err := r.db.Query(ctx, q, orgID, since, limit)
if err != nil {
return nil, err
}
defer rows.Close()
var out []ThreadFollowUpState
for rows.Next() {
var st ThreadFollowUpState
var lastIn, lastOut *time.Time
if err := rows.Scan(&st.ThreadID, &lastIn, &lastOut, &st.BestIntent, &st.LastKind); err != nil {
return nil, err
}
if lastIn != nil {
st.LastInboundAt = *lastIn
}
if lastOut != nil {
st.LastOutboundAt = *lastOut
}
out = append(out, st)
}
return out, rows.Err()
}