Skip to content

Commit d93adad

Browse files
authored
feat(push_delivery,backend): cut over push delivery and RSS handler to BaseSummaryVisitor (#2622) (#2637)
* feat(push_delivery,backend): cut over push delivery and RSS handler to BaseSummaryVisitor (#2622) Complete the Visitor pattern cutover across push delivery pipelines and the RSS HTTP handler. - Update shouldNotifyV1 in workers/push_delivery/pkg/dispatcher to delegate highlight filtering and error checks to EventSummary.Categorize(triggers) and BaseSummaryVisitor.HasContent(). - Cut over get_subscription_rss.go in backend/pkg/httpserver to invoke summary.Accept(visitor, triggers) for RSS feed payload generation. - Pass EventSummary by pointer (*EventSummary) into shouldNotifyV1 to avoid header and slice backing array copy overhead across subscriber iteration loops. - Remove redundant manual trigger matching logic, completing the double-dispatch cutover for push notifications and RSS feed rendering. CONV=f70d8c59-6f49-4a69-bb74-5643e908f5b1 TAG=agy * address feedback
1 parent 3ebce71 commit d93adad

3 files changed

Lines changed: 115 additions & 125 deletions

File tree

backend/pkg/httpserver/get_subscription_rss.go

Lines changed: 29 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@ package httpserver
1717
import (
1818
"bytes"
1919
"context"
20+
"encoding/json"
2021
"encoding/xml"
2122
"errors"
2223
"fmt"
@@ -176,22 +177,11 @@ func (s *Server) GetSubscriptionRSS(
176177
})
177178
}
178179

179-
var jobTriggers []workertypes.JobTrigger
180-
for _, triggerItem := range sub.Triggers {
181-
triggerVal, err := triggerItem.Value.AsSubscriptionTriggerWritable()
182-
if err != nil {
183-
continue
184-
}
185-
if jobTrigger, ok := workertypes.ToJobTrigger(triggerVal); ok {
186-
jobTriggers = append(jobTriggers, jobTrigger)
187-
}
188-
}
180+
jobTriggers := extractJobTriggers(sub.Triggers)
189181

190182
for _, e := range events {
191-
visitor := newRSSVisitor(jobTriggers)
192-
var description string
193-
var title string
194-
if err := workertypes.ParseEventSummary(e.Summary, visitor); err != nil {
183+
var summary workertypes.EventSummary
184+
if err := json.Unmarshal([]byte(e.Summary), &summary); err != nil {
195185
slog.ErrorContext(ctx, "failed to unmarshal summary", "event_id", e.ID, "error", err)
196186

197187
errorHTML := fmt.Sprintf(
@@ -212,6 +202,16 @@ func (s *Server) GetSubscriptionRSS(
212202
continue
213203
}
214204

205+
var description string
206+
var title string
207+
208+
visitor := newRSSVisitor(jobTriggers)
209+
if err := visitor.VisitV1(summary); err != nil {
210+
slog.ErrorContext(ctx, "failed to process RSS summary visitor", "event_id", e.ID, "error", err)
211+
212+
continue
213+
}
214+
215215
if !visitor.HasContent() {
216216
continue
217217
}
@@ -262,3 +262,18 @@ func (s *Server) GetSubscriptionRSS(
262262
ContentLength: int64(buf.Len()),
263263
}, nil
264264
}
265+
266+
func extractJobTriggers(triggers []backend.SubscriptionTriggerResponseItem) []workertypes.JobTrigger {
267+
var jobTriggers []workertypes.JobTrigger
268+
for _, triggerItem := range triggers {
269+
triggerVal, err := triggerItem.Value.AsSubscriptionTriggerWritable()
270+
if err != nil {
271+
continue
272+
}
273+
if jobTrigger, ok := workertypes.ToJobTrigger(triggerVal); ok {
274+
jobTriggers = append(jobTriggers, jobTrigger)
275+
}
276+
}
277+
278+
return jobTriggers
279+
}

workers/push_delivery/pkg/dispatcher/dispatcher.go

Lines changed: 22 additions & 35 deletions
Original file line numberDiff line numberDiff line change
@@ -159,7 +159,11 @@ func (g *deliveryJobGenerator) VisitV1(s workertypes.EventSummary) error {
159159
// 2. Filter & Create Jobs
160160
// Iterate Emails.
161161
for _, sub := range subscribers.Emails {
162-
if !shouldNotifyV1(sub.Triggers, s) {
162+
notify, err := shouldNotifyV1(sub.Triggers, &s)
163+
if err != nil {
164+
return fmt.Errorf("error checking notification triggers for email subscription %s: %w", sub.SubscriptionID, err)
165+
}
166+
if !notify {
163167
continue
164168
}
165169
g.emailJobs = append(g.emailJobs, workertypes.EmailDeliveryJob{
@@ -174,7 +178,11 @@ func (g *deliveryJobGenerator) VisitV1(s workertypes.EventSummary) error {
174178

175179
// Iterate Webhooks.
176180
for _, sub := range subscribers.Webhooks {
177-
if !shouldNotifyV1(sub.Triggers, s) {
181+
notify, err := shouldNotifyV1(sub.Triggers, &s)
182+
if err != nil {
183+
return fmt.Errorf("error checking notification triggers for webhook subscription %s: %w", sub.SubscriptionID, err)
184+
}
185+
if !notify {
178186
continue
179187
}
180188
g.webhookJobs = append(g.webhookJobs, workertypes.WebhookDeliveryJob{
@@ -196,40 +204,19 @@ func (g *deliveryJobGenerator) JobCount() int {
196204
return len(g.emailJobs) + len(g.webhookJobs)
197205
}
198206

199-
// shouldNotifyV1 determines if the V1 event summary matches any of the user's triggers.
200-
func shouldNotifyV1(triggers []workertypes.JobTrigger, summary workertypes.EventSummary) bool {
201-
if len(summary.QueryErrors) > 0 || len(summary.ResolvedQueryErrors) > 0 {
202-
return true
207+
// shouldNotifyV1 determines if the V1 event summary matches any of the subscriber's triggers.
208+
// Note: summary is passed as a pointer (*workertypes.EventSummary) to prevent copying the
209+
// EventSummary struct header and slice backing arrays across subscriber evaluation loops.
210+
// base.HasContent() returns true if there are active query errors, resolved query errors, or any
211+
// highlights matching the given triggers.
212+
func shouldNotifyV1(triggers []workertypes.JobTrigger, summary *workertypes.EventSummary) (bool, error) {
213+
if summary == nil {
214+
return false, nil
203215
}
204-
205-
// 1. Determine if summary has changes.
206-
hasChanges := summary.Categories.Added > 0 ||
207-
summary.Categories.Removed > 0 ||
208-
summary.Categories.Updated > 0 ||
209-
summary.Categories.Moved > 0 ||
210-
summary.Categories.Split > 0 ||
211-
summary.Categories.QueryChanged > 0
212-
213-
if !hasChanges {
214-
return false
215-
}
216-
217-
// 2. Iterate triggers and check highlights.
218-
for _, t := range triggers {
219-
if matchesTrigger(t, summary) {
220-
return true
221-
}
222-
}
223-
224-
return false
225-
}
226-
227-
func matchesTrigger(t workertypes.JobTrigger, summary workertypes.EventSummary) bool {
228-
for _, h := range summary.Highlights {
229-
if h.MatchesTrigger(t) {
230-
return true
231-
}
216+
base, err := summary.Categorize(triggers)
217+
if err != nil {
218+
return false, fmt.Errorf("failed to categorize event summary against triggers: %w", err)
232219
}
233220

234-
return false
221+
return base.HasContent(), nil
235222
}

workers/push_delivery/pkg/dispatcher/dispatcher_test.go

Lines changed: 64 additions & 76 deletions
Original file line numberDiff line numberDiff line change
@@ -97,42 +97,21 @@ func createTestSummary(hasChanges bool) workertypes.EventSummary {
9797
categories.Added = 1
9898
}
9999

100-
return workertypes.EventSummary{
101-
SchemaVersion: "v1",
102-
SnapshotOrigin: workertypes.OriginLive,
103-
Text: "Test Summary",
104-
Categories: categories,
105-
Truncated: false,
106-
QueryErrors: nil,
107-
ResolvedQueryErrors: nil,
108-
Highlights: nil,
109-
}
100+
summary := workertypes.NewEmptyEventSummary()
101+
summary.SnapshotOrigin = workertypes.OriginLive
102+
summary.Text = "Test Summary"
103+
summary.Categories = categories
104+
105+
return summary
110106
}
111107

112108
func createTestSummaryWithErrors(errCode workertypes.SummaryQueryErrorCode) workertypes.EventSummary {
113-
categories := workertypes.SummaryCategories{
114-
QueryChanged: 0,
115-
Added: 0,
116-
Deleted: 0,
117-
Removed: 0,
118-
Moved: 0,
119-
Split: 0,
120-
Updated: 0,
121-
UpdatedImpl: 0,
122-
UpdatedRename: 0,
123-
UpdatedBaseline: 0,
124-
}
109+
summary := workertypes.NewEmptyEventSummary()
110+
summary.SnapshotOrigin = workertypes.OriginLive
111+
summary.Text = "Error occurred"
112+
summary.SetQueryErrors([]workertypes.SummaryQueryError{{Code: errCode}})
125113

126-
return workertypes.EventSummary{
127-
SchemaVersion: workertypes.VersionEventSummaryV1,
128-
SnapshotOrigin: workertypes.OriginLive,
129-
Text: "Error occurred",
130-
Categories: categories,
131-
Truncated: false,
132-
QueryErrors: []workertypes.SummaryQueryError{{Code: errCode}},
133-
ResolvedQueryErrors: nil,
134-
Highlights: nil,
135-
}
114+
return summary
136115
}
137116

138117
// mockParserFactory creates a SummaryParser that injects the given summary directly.
@@ -205,22 +184,20 @@ func TestProcessEvent_Success(t *testing.T) {
205184
summary := createTestSummary(true)
206185
summary.Categories.UpdatedBaseline = 1
207186
summary.Categories.Updated = 1
208-
summary.Highlights = []workertypes.SummaryHighlight{
209-
{
210-
Type: workertypes.SummaryHighlightTypeChanged,
211-
FeatureID: "test-feature-id",
212-
FeatureName: "Test Feature",
213-
Docs: nil,
214-
NameChange: nil,
215-
BaselineChange: &workertypes.Change[workertypes.BaselineValue]{
216-
From: newBaselineValue(workertypes.BaselineStatusLimited),
217-
To: newBaselineValue(workertypes.BaselineStatusNewly),
218-
},
219-
BrowserChanges: nil,
220-
Moved: nil,
221-
Split: nil,
187+
summary.AddHighlight(workertypes.SummaryHighlight{
188+
Type: workertypes.SummaryHighlightTypeChanged,
189+
FeatureID: "test-feature-id",
190+
FeatureName: "Test Feature",
191+
Docs: nil,
192+
NameChange: nil,
193+
BaselineChange: &workertypes.Change[workertypes.BaselineValue]{
194+
From: newBaselineValue(workertypes.BaselineStatusLimited),
195+
To: newBaselineValue(workertypes.BaselineStatusNewly),
222196
},
223-
}
197+
BrowserChanges: nil,
198+
Moved: nil,
199+
Split: nil,
200+
})
224201
parser := mockParserFactory(summary, nil)
225202

226203
d := NewDispatcher(finder, publisher)
@@ -315,22 +292,20 @@ func TestProcessEvent_Webhook_Success(t *testing.T) {
315292
summary := createTestSummary(true)
316293
summary.Categories.UpdatedBaseline = 1
317294
summary.Categories.Updated = 1
318-
summary.Highlights = []workertypes.SummaryHighlight{
319-
{
320-
Type: workertypes.SummaryHighlightTypeChanged,
321-
FeatureID: "test-feature-id",
322-
FeatureName: "Test Feature",
323-
Docs: nil,
324-
NameChange: nil,
325-
BaselineChange: &workertypes.Change[workertypes.BaselineValue]{
326-
From: newBaselineValue(workertypes.BaselineStatusLimited),
327-
To: newBaselineValue(workertypes.BaselineStatusNewly),
328-
},
329-
BrowserChanges: nil,
330-
Moved: nil,
331-
Split: nil,
295+
summary.AddHighlight(workertypes.SummaryHighlight{
296+
Type: workertypes.SummaryHighlightTypeChanged,
297+
FeatureID: "test-feature-id",
298+
FeatureName: "Test Feature",
299+
Docs: nil,
300+
NameChange: nil,
301+
BaselineChange: &workertypes.Change[workertypes.BaselineValue]{
302+
From: newBaselineValue(workertypes.BaselineStatusLimited),
303+
To: newBaselineValue(workertypes.BaselineStatusNewly),
332304
},
333-
}
305+
BrowserChanges: nil,
306+
Moved: nil,
307+
Split: nil,
308+
})
334309
parser := mockParserFactory(summary, nil)
335310

336311
d := NewDispatcher(finder, publisher)
@@ -601,7 +576,7 @@ func newBrowserValue(status workertypes.BrowserStatus) workertypes.BrowserValue
601576

602577
func withBaselineHighlight(
603578
s workertypes.EventSummary, from, to workertypes.BaselineStatus) workertypes.EventSummary {
604-
s.Highlights = append(s.Highlights, workertypes.SummaryHighlight{
579+
s.AddHighlight(workertypes.SummaryHighlight{
605580
Type: workertypes.SummaryHighlightTypeChanged,
606581
FeatureID: "test-feature-id",
607582
FeatureName: "Test Feature",
@@ -623,7 +598,7 @@ func withBaselineHighlight(
623598

624599
func withBrowserChangeHighlight(
625600
s workertypes.EventSummary, from, to workertypes.BrowserStatus) workertypes.EventSummary {
626-
s.Highlights = append(s.Highlights, workertypes.SummaryHighlight{
601+
s.AddHighlight(workertypes.SummaryHighlight{
627602
Type: workertypes.SummaryHighlightTypeChanged,
628603
FeatureID: "test-feature-id",
629604
FeatureName: "Test Feature",
@@ -764,7 +739,10 @@ func TestShouldNotifyV1(t *testing.T) {
764739

765740
for _, tc := range testCases {
766741
t.Run(tc.name, func(t *testing.T) {
767-
got := shouldNotifyV1(tc.triggers, tc.summary)
742+
got, err := shouldNotifyV1(tc.triggers, &tc.summary)
743+
if err != nil {
744+
t.Fatalf("shouldNotifyV1 unexpected error: %v", err)
745+
}
768746
if got != tc.want {
769747
t.Errorf("shouldNotifyV1() = %v, want %v", got, tc.want)
770748
}
@@ -773,16 +751,12 @@ func TestShouldNotifyV1(t *testing.T) {
773751
}
774752

775753
func createRecoveredSummary() workertypes.EventSummary {
776-
return workertypes.EventSummary{
777-
SchemaVersion: workertypes.VersionEventSummaryV1,
778-
SnapshotOrigin: workertypes.OriginLive,
779-
Text: "Search query recovered and tracking 2 features normally.",
780-
Truncated: false,
781-
Categories: workertypes.NewEmptySummaryCategories(),
782-
QueryErrors: nil,
783-
ResolvedQueryErrors: []workertypes.SummaryQueryError{{Code: workertypes.SummaryQueryErrorCodeQueryGrammar}},
784-
Highlights: nil,
785-
}
754+
summary := workertypes.NewEmptyEventSummary()
755+
summary.SnapshotOrigin = workertypes.OriginLive
756+
summary.Text = "Search query recovered and tracking 2 features normally."
757+
summary.SetResolvedQueryErrors([]workertypes.SummaryQueryError{{Code: workertypes.SummaryQueryErrorCodeQueryGrammar}})
758+
759+
return summary
786760
}
787761

788762
func TestShouldNotifyV1_ResolvedQueryErrors(t *testing.T) {
@@ -816,7 +790,10 @@ func TestShouldNotifyV1_ResolvedQueryErrors(t *testing.T) {
816790

817791
for _, tc := range testCases {
818792
t.Run(tc.name, func(t *testing.T) {
819-
got := shouldNotifyV1(tc.triggers, tc.summary)
793+
got, err := shouldNotifyV1(tc.triggers, &tc.summary)
794+
if err != nil {
795+
t.Fatalf("shouldNotifyV1 unexpected error: %v", err)
796+
}
820797
if got != tc.want {
821798
t.Errorf("shouldNotifyV1() = %v, want %v", got, tc.want)
822799
}
@@ -878,3 +855,14 @@ func TestProcessEvent_ResolvedQueryErrors_AllChannels(t *testing.T) {
878855
t.Errorf("expected subscription ID sub-webhook-recovery, got %s", publisher.webhookJobs[0].SubscriptionID)
879856
}
880857
}
858+
859+
func TestShouldNotifyV1_NilSummary(t *testing.T) {
860+
// Verify that if the summary is nil, shouldNotifyV1 returns false and no error.
861+
got, err := shouldNotifyV1([]workertypes.JobTrigger{workertypes.FeaturePromotedToNewly}, nil)
862+
if err != nil {
863+
t.Fatalf("shouldNotifyV1 unexpected error: %v", err)
864+
}
865+
if got {
866+
t.Errorf("shouldNotifyV1(triggers, nil) = %v, want false", got)
867+
}
868+
}

0 commit comments

Comments
 (0)