diff --git a/cmd/warmblyctl/inboxtag.go b/cmd/warmblyctl/inboxtag.go index 94ed81721..8961ef3ec 100644 --- a/cmd/warmblyctl/inboxtag.go +++ b/cmd/warmblyctl/inboxtag.go @@ -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) } diff --git a/docs/content/docs/guides/inbox-tagging.mdx b/docs/content/docs/guides/inbox-tagging.mdx index 0942c4775..134040f19 100644 --- a/docs/content/docs/guides/inbox-tagging.mdx +++ b/docs/content/docs/guides/inbox-tagging.mdx @@ -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 diff --git a/internal/app/consumer/service.go b/internal/app/consumer/service.go index 84da2ea0e..263bb81a4 100644 --- a/internal/app/consumer/service.go +++ b/internal/app/consumer/service.go @@ -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()}) } diff --git a/internal/app/inboxtag/custom_test.go b/internal/app/inboxtag/custom_test.go index 72e9600ed..6c9c58458 100644 --- a/internal/app/inboxtag/custom_test.go +++ b/internal/app/inboxtag/custom_test.go @@ -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"} { diff --git a/internal/app/inboxtag/followup_test.go b/internal/app/inboxtag/followup_test.go index a7b7f1027..f33d7f27a 100644 --- a/internal/app/inboxtag/followup_test.go +++ b/internal/app/inboxtag/followup_test.go @@ -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) + } +} diff --git a/internal/app/inboxtag/inboxtag_test.go b/internal/app/inboxtag/inboxtag_test.go index c6682401e..0115aad59 100644 --- a/internal/app/inboxtag/inboxtag_test.go +++ b/internal/app/inboxtag/inboxtag_test.go @@ -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) { diff --git a/internal/app/inboxtag/service.go b/internal/app/inboxtag/service.go index 6811c460c..127b52638 100644 --- a/internal/app/inboxtag/service.go +++ b/internal/app/inboxtag/service.go @@ -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 -} diff --git a/internal/app/inboxtag/sweep.go b/internal/app/inboxtag/sweep.go new file mode 100644 index 000000000..6bd53b475 --- /dev/null +++ b/internal/app/inboxtag/sweep.go @@ -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 +} diff --git a/internal/app/inboxtag/sweep_live_test.go b/internal/app/inboxtag/sweep_live_test.go new file mode 100644 index 000000000..3cd16e585 --- /dev/null +++ b/internal/app/inboxtag/sweep_live_test.go @@ -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/?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, 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 := "" + 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', '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) + } +} diff --git a/internal/app/orgtransfer/spec.go b/internal/app/orgtransfer/spec.go index d88d09163..5b703f110 100644 --- a/internal/app/orgtransfer/spec.go +++ b/internal/app/orgtransfer/spec.go @@ -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.", } diff --git a/internal/infrastructure/db/migrations/000245_inbox_follow_up_sweeps.down.sql b/internal/infrastructure/db/migrations/000245_inbox_follow_up_sweeps.down.sql new file mode 100644 index 000000000..51c175475 --- /dev/null +++ b/internal/infrastructure/db/migrations/000245_inbox_follow_up_sweeps.down.sql @@ -0,0 +1 @@ +DROP TABLE IF EXISTS public.inbox_follow_up_sweeps; diff --git a/internal/infrastructure/db/migrations/000245_inbox_follow_up_sweeps.up.sql b/internal/infrastructure/db/migrations/000245_inbox_follow_up_sweeps.up.sql new file mode 100644 index 000000000..565486f13 --- /dev/null +++ b/internal/infrastructure/db/migrations/000245_inbox_follow_up_sweeps.up.sql @@ -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() +); diff --git a/internal/infrastructure/db/migrations/000246_unibox_emails_account_date_index.down.sql b/internal/infrastructure/db/migrations/000246_unibox_emails_account_date_index.down.sql new file mode 100644 index 000000000..72601a12f --- /dev/null +++ b/internal/infrastructure/db/migrations/000246_unibox_emails_account_date_index.down.sql @@ -0,0 +1 @@ +DROP INDEX CONCURRENTLY IF EXISTS idx_unibox_emails_account_date; diff --git a/internal/infrastructure/db/migrations/000246_unibox_emails_account_date_index.up.sql b/internal/infrastructure/db/migrations/000246_unibox_emails_account_date_index.up.sql new file mode 100644 index 000000000..bfbf33fd2 --- /dev/null +++ b/internal/infrastructure/db/migrations/000246_unibox_emails_account_date_index.up.sql @@ -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); diff --git a/internal/infrastructure/db/migrations/000247_unibox_emails_ingested_at.down.sql b/internal/infrastructure/db/migrations/000247_unibox_emails_ingested_at.down.sql new file mode 100644 index 000000000..c45038c15 --- /dev/null +++ b/internal/infrastructure/db/migrations/000247_unibox_emails_ingested_at.down.sql @@ -0,0 +1 @@ +ALTER TABLE public.unibox_emails DROP COLUMN IF EXISTS ingested_at; diff --git a/internal/infrastructure/db/migrations/000247_unibox_emails_ingested_at.up.sql b/internal/infrastructure/db/migrations/000247_unibox_emails_ingested_at.up.sql new file mode 100644 index 000000000..e2b92149a --- /dev/null +++ b/internal/infrastructure/db/migrations/000247_unibox_emails_ingested_at.up.sql @@ -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(); diff --git a/internal/infrastructure/db/migrations/000248_unibox_emails_account_ingested_index.down.sql b/internal/infrastructure/db/migrations/000248_unibox_emails_account_ingested_index.down.sql new file mode 100644 index 000000000..510bc941f --- /dev/null +++ b/internal/infrastructure/db/migrations/000248_unibox_emails_account_ingested_index.down.sql @@ -0,0 +1 @@ +DROP INDEX CONCURRENTLY IF EXISTS idx_unibox_emails_account_ingested; diff --git a/internal/infrastructure/db/migrations/000248_unibox_emails_account_ingested_index.up.sql b/internal/infrastructure/db/migrations/000248_unibox_emails_account_ingested_index.up.sql new file mode 100644 index 000000000..6f8af9f29 --- /dev/null +++ b/internal/infrastructure/db/migrations/000248_unibox_emails_account_ingested_index.up.sql @@ -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); diff --git a/internal/infrastructure/db/migrations/000249_inbox_tag_results_updated_index.down.sql b/internal/infrastructure/db/migrations/000249_inbox_tag_results_updated_index.down.sql new file mode 100644 index 000000000..8945626ca --- /dev/null +++ b/internal/infrastructure/db/migrations/000249_inbox_tag_results_updated_index.down.sql @@ -0,0 +1 @@ +DROP INDEX CONCURRENTLY IF EXISTS idx_inbox_tag_results_updated; diff --git a/internal/infrastructure/db/migrations/000249_inbox_tag_results_updated_index.up.sql b/internal/infrastructure/db/migrations/000249_inbox_tag_results_updated_index.up.sql new file mode 100644 index 000000000..229821901 --- /dev/null +++ b/internal/infrastructure/db/migrations/000249_inbox_tag_results_updated_index.up.sql @@ -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); diff --git a/internal/repository/pg_inbox_follow_up.go b/internal/repository/pg_inbox_follow_up.go new file mode 100644 index 000000000..df29b9564 --- /dev/null +++ b/internal/repository/pg_inbox_follow_up.go @@ -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 +} diff --git a/internal/repository/pg_inbox_tag.go b/internal/repository/pg_inbox_tag.go index b6fe72303..10631a007 100644 --- a/internal/repository/pg_inbox_tag.go +++ b/internal/repository/pg_inbox_tag.go @@ -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() -}