Skip to content

Commit 6a3222c

Browse files
fix: publish segment domain events after commit to prevent stale cache refresh
Co-authored-by: Cursor <cursoragent@cursor.com>
1 parent 088ffbc commit 6a3222c

4 files changed

Lines changed: 546 additions & 39 deletions

File tree

pkg/feature/api/segment.go

Lines changed: 72 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -94,24 +94,7 @@ func (s *FeatureService) CreateSegment(
9494
)
9595
return err
9696
}
97-
e, err := domainevent.NewEvent(
98-
editor,
99-
eventproto.Event_SEGMENT,
100-
segment.Id,
101-
eventproto.Event_SEGMENT_CREATED,
102-
&eventproto.SegmentCreatedEvent{
103-
Id: segment.Id,
104-
Name: segment.Name,
105-
Description: segment.Description,
106-
},
107-
req.EnvironmentId,
108-
segment.Segment,
109-
nil,
110-
)
111-
if err != nil {
112-
return nil
113-
}
114-
return s.domainPublisher.Publish(ctx, e)
97+
return nil
11598
})
11699
if err != nil {
117100
if errors.Is(err, v2fs.ErrSegmentAlreadyExists) {
@@ -126,6 +109,44 @@ func (s *FeatureService) CreateSegment(
126109
)
127110
return nil, api.NewGRPCStatus(err).Err()
128111
}
112+
// Publish only after the transaction commits: consumers such as the
113+
// cache refresher re-read MySQL on each event, so publishing inside
114+
// the transaction would let them read (and cache) pre-commit state.
115+
e, err := domainevent.NewEvent(
116+
editor,
117+
eventproto.Event_SEGMENT,
118+
segment.Id,
119+
eventproto.Event_SEGMENT_CREATED,
120+
&eventproto.SegmentCreatedEvent{
121+
Id: segment.Id,
122+
Name: segment.Name,
123+
Description: segment.Description,
124+
},
125+
req.EnvironmentId,
126+
segment.Segment,
127+
nil,
128+
)
129+
if err != nil {
130+
s.logger.Error(
131+
"Failed to create domain event",
132+
log.FieldsFromIncomingContext(ctx).AddFields(
133+
zap.Error(err),
134+
zap.String("environmentId", req.EnvironmentId),
135+
)...,
136+
)
137+
return nil, api.NewGRPCStatus(err).Err()
138+
}
139+
if err := s.domainPublisher.Publish(ctx, e); err != nil {
140+
s.logger.Error(
141+
"Failed to publish domain event",
142+
log.FieldsFromIncomingContext(ctx).AddFields(
143+
zap.Error(err),
144+
zap.String("environmentId", req.EnvironmentId),
145+
zap.Any("event", e),
146+
)...,
147+
)
148+
return nil, api.NewGRPCStatus(err).Err()
149+
}
129150
return &featureproto.CreateSegmentResponse{
130151
Segment: segment.Segment,
131152
}, nil
@@ -154,6 +175,7 @@ func (s *FeatureService) DeleteSegment(
154175
if err := s.checkSegmentInUse(ctx, req.Id, req.EnvironmentId); err != nil {
155176
return nil, err
156177
}
178+
var eventPb *eventproto.Event
157179
err = s.dbClient.RunInTransactionV2(ctx, func(contextWithTx context.Context) error {
158180
segment, _, err := s.segmentStorage.GetSegment(contextWithTx, req.Id, req.EnvironmentId)
159181
if err != nil {
@@ -166,7 +188,7 @@ func (s *FeatureService) DeleteSegment(
166188
)
167189
return err
168190
}
169-
event, err := domainevent.NewEvent(
191+
eventPb, err = domainevent.NewEvent(
170192
editor,
171193
eventproto.Event_SEGMENT,
172194
segment.Id,
@@ -179,9 +201,6 @@ func (s *FeatureService) DeleteSegment(
179201
segment.Segment, // Previous state: what was deleted
180202
)
181203
if err != nil {
182-
return nil
183-
}
184-
if err := s.domainPublisher.Publish(ctx, event); err != nil {
185204
return err
186205
}
187206
return s.segmentStorage.DeleteSegment(contextWithTx, segment.Id)
@@ -198,6 +217,20 @@ func (s *FeatureService) DeleteSegment(
198217
)
199218
return nil, api.NewGRPCStatus(err).Err()
200219
}
220+
// Publish only after the transaction commits: consumers such as the
221+
// cache refresher act on the event immediately, so publishing inside
222+
// the transaction would let them observe pre-commit state.
223+
if err := s.domainPublisher.Publish(ctx, eventPb); err != nil {
224+
s.logger.Error(
225+
"Failed to publish domain event",
226+
log.FieldsFromIncomingContext(ctx).AddFields(
227+
zap.Error(err),
228+
zap.String("environmentId", req.EnvironmentId),
229+
zap.Any("event", eventPb),
230+
)...,
231+
)
232+
return nil, api.NewGRPCStatus(err).Err()
233+
}
201234
return &featureproto.DeleteSegmentResponse{}, nil
202235
}
203236

@@ -291,6 +324,7 @@ func (s *FeatureService) UpdateSegment(
291324
return nil, err
292325
}
293326
var updatedSegment *featureproto.Segment
327+
var eventPb *eventproto.Event
294328
err = s.dbClient.RunInTransactionV2(ctx, func(contextWithTx context.Context) error {
295329
segment, _, err := s.segmentStorage.GetSegment(contextWithTx, req.Id, req.EnvironmentId)
296330
if err != nil {
@@ -319,7 +353,7 @@ func (s *FeatureService) UpdateSegment(
319353
updated.UpdateRules(req.Rules.Values)
320354
}
321355
updatedSegment = updated.Segment
322-
e, err := domainevent.NewEvent(
356+
eventPb, err = domainevent.NewEvent(
323357
editor,
324358
eventproto.Event_SEGMENT,
325359
req.Id,
@@ -337,9 +371,6 @@ func (s *FeatureService) UpdateSegment(
337371
if err != nil {
338372
return err
339373
}
340-
if err := s.domainPublisher.Publish(ctx, e); err != nil {
341-
return err
342-
}
343374
return s.segmentStorage.UpdateSegment(contextWithTx, updated, req.EnvironmentId)
344375
})
345376
if err != nil {
@@ -355,6 +386,21 @@ func (s *FeatureService) UpdateSegment(
355386
)
356387
return nil, api.NewGRPCStatus(err).Err()
357388
}
389+
// Publish only after the transaction commits: the cache refresher
390+
// re-reads MySQL on each event, so publishing inside the transaction
391+
// would let it read the pre-update segment and overwrite the segment
392+
// users cache with stale rules.
393+
if err := s.domainPublisher.Publish(ctx, eventPb); err != nil {
394+
s.logger.Error(
395+
"Failed to publish domain event",
396+
log.FieldsFromIncomingContext(ctx).AddFields(
397+
zap.Error(err),
398+
zap.String("environmentId", req.EnvironmentId),
399+
zap.Any("event", eventPb),
400+
)...,
401+
)
402+
return nil, api.NewGRPCStatus(err).Err()
403+
}
358404
// Refresh the segment users cache so evaluation and server SDK sync paths
359405
// pick up the rule change without waiting for the batch cacher.
360406
// The cache refresh is best effort: on failure the batch cacher will

pkg/feature/api/segment_test.go

Lines changed: 208 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@ package api
1616

1717
import (
1818
"context"
19+
"errors"
1920
"testing"
2021
"time"
2122

@@ -31,6 +32,8 @@ import (
3132
"github.com/bucketeer-io/bucketeer/v2/pkg/feature/domain"
3233
v2fs "github.com/bucketeer-io/bucketeer/v2/pkg/feature/storage/v2"
3334
storagemock "github.com/bucketeer-io/bucketeer/v2/pkg/feature/storage/v2/mock"
35+
"github.com/bucketeer-io/bucketeer/v2/pkg/pubsub/publisher"
36+
publishermock "github.com/bucketeer-io/bucketeer/v2/pkg/pubsub/publisher/mock"
3437
"github.com/bucketeer-io/bucketeer/v2/pkg/rpc"
3538
databasemock "github.com/bucketeer-io/bucketeer/v2/pkg/storage/v2/database/mock"
3639
"github.com/bucketeer-io/bucketeer/v2/pkg/token"
@@ -732,6 +735,211 @@ func TestListSegmentsMySQL(t *testing.T) {
732735
}
733736
}
734737

738+
// TestSegmentDomainEventPublishedAfterCommitMySQL guards the
739+
// publish-after-commit ordering of segment mutations. Publishing inside the
740+
// transaction lets consumers that re-read MySQL on each event (e.g. the
741+
// cache refresher) observe pre-commit state and overwrite the segment users
742+
// cache with stale rules, which broke server SDK diff syncs.
743+
func TestSegmentDomainEventPublishedAfterCommitMySQL(t *testing.T) {
744+
t.Parallel()
745+
mockController := gomock.NewController(t)
746+
defer mockController.Finish()
747+
748+
ctx, cancel := context.WithCancel(context.Background())
749+
defer cancel()
750+
ctx = metadata.NewIncomingContext(ctx, metadata.MD{
751+
"accept-language": []string{"ja"},
752+
})
753+
ctx = setToken(ctx)
754+
755+
segmentRules := &featureproto.RuleListValue{
756+
Values: []*featureproto.Rule{
757+
{
758+
Clauses: []*featureproto.Clause{
759+
{
760+
Attribute: "plan",
761+
Operator: featureproto.Clause_EQUALS,
762+
Values: []string{"premium"},
763+
},
764+
},
765+
},
766+
},
767+
}
768+
769+
testcases := []struct {
770+
desc string
771+
setup func(s *FeatureService, committed *bool)
772+
run func(s *FeatureService) error
773+
}{
774+
{
775+
desc: "CreateSegment",
776+
setup: func(s *FeatureService, committed *bool) {
777+
s.dbClient.(*databasemock.MockClient).EXPECT().RunInTransactionV2(
778+
gomock.Any(), gomock.Any(),
779+
).DoAndReturn(func(ctx context.Context, fn func(ctx context.Context) error) error {
780+
err := fn(ctx)
781+
require.NoError(t, err)
782+
*committed = true
783+
return err
784+
})
785+
s.segmentStorage.(*storagemock.MockSegmentStorage).EXPECT().CreateSegment(
786+
gomock.Any(), gomock.Any(), gomock.Any(),
787+
).Return(nil)
788+
},
789+
run: func(s *FeatureService) error {
790+
_, err := s.CreateSegment(ctx, &featureproto.CreateSegmentRequest{
791+
Name: "name",
792+
Description: "description",
793+
EnvironmentId: "ns0",
794+
})
795+
return err
796+
},
797+
},
798+
{
799+
desc: "UpdateSegment with rules",
800+
setup: func(s *FeatureService, committed *bool) {
801+
s.dbClient.(*databasemock.MockClient).EXPECT().RunInTransactionV2(
802+
gomock.Any(), gomock.Any(),
803+
).DoAndReturn(func(ctx context.Context, fn func(ctx context.Context) error) error {
804+
err := fn(ctx)
805+
require.NoError(t, err)
806+
*committed = true
807+
return err
808+
})
809+
s.segmentStorage.(*storagemock.MockSegmentStorage).EXPECT().GetSegment(
810+
gomock.Any(), gomock.Any(), gomock.Any(),
811+
).Return(&domain.Segment{
812+
Segment: &featureproto.Segment{
813+
Id: "id0",
814+
},
815+
}, nil, nil)
816+
s.segmentStorage.(*storagemock.MockSegmentStorage).EXPECT().UpdateSegment(
817+
gomock.Any(), gomock.Any(), gomock.Any(),
818+
).Return(nil)
819+
s.segmentStorage.(*storagemock.MockSegmentStorage).EXPECT().ListSegmentUsersBySegment(
820+
gomock.Any(), "id0", "ns0",
821+
).Return([]*featureproto.SegmentUser{}, nil)
822+
s.segmentUsersCache.(*cachev3mock.MockSegmentUsersCache).EXPECT().Put(
823+
gomock.Any(), "ns0",
824+
).Return(nil)
825+
},
826+
run: func(s *FeatureService) error {
827+
_, err := s.UpdateSegment(ctx, &featureproto.UpdateSegmentRequest{
828+
Id: "id0",
829+
EnvironmentId: "ns0",
830+
Rules: segmentRules,
831+
})
832+
return err
833+
},
834+
},
835+
{
836+
desc: "DeleteSegment",
837+
setup: func(s *FeatureService, committed *bool) {
838+
s.featureStorage.(*storagemock.MockFeatureStorage).EXPECT().ListFeatures(
839+
gomock.Any(), gomock.Any(),
840+
).Return([]*featureproto.Feature{}, 0, int64(0), nil)
841+
s.dbClient.(*databasemock.MockClient).EXPECT().RunInTransactionV2(
842+
gomock.Any(), gomock.Any(),
843+
).DoAndReturn(func(ctx context.Context, fn func(ctx context.Context) error) error {
844+
err := fn(ctx)
845+
require.NoError(t, err)
846+
*committed = true
847+
return err
848+
})
849+
s.segmentStorage.(*storagemock.MockSegmentStorage).EXPECT().GetSegment(
850+
gomock.Any(), gomock.Any(), gomock.Any(),
851+
).Return(&domain.Segment{
852+
Segment: &featureproto.Segment{
853+
Id: "id0",
854+
},
855+
}, nil, nil)
856+
s.segmentStorage.(*storagemock.MockSegmentStorage).EXPECT().DeleteSegment(
857+
gomock.Any(), gomock.Any(),
858+
).Return(nil)
859+
},
860+
run: func(s *FeatureService) error {
861+
_, err := s.DeleteSegment(ctx, &featureproto.DeleteSegmentRequest{
862+
Id: "id0",
863+
EnvironmentId: "ns0",
864+
})
865+
return err
866+
},
867+
},
868+
}
869+
for _, tc := range testcases {
870+
t.Run(tc.desc, func(t *testing.T) {
871+
service := createFeatureService(mockController)
872+
domainPublisher := publishermock.NewMockPublisher(mockController)
873+
service.domainPublisher = domainPublisher
874+
committed := false
875+
domainPublisher.EXPECT().Publish(gomock.Any(), gomock.Any()).DoAndReturn(
876+
func(ctx context.Context, msg publisher.Message) error {
877+
assert.True(t, committed,
878+
"domain event must be published after the transaction commits")
879+
return nil
880+
})
881+
tc.setup(service, &committed)
882+
assert.NoError(t, tc.run(service))
883+
})
884+
}
885+
}
886+
887+
// TestUpdateSegmentPublishFailureMySQL: when the post-commit publish fails,
888+
// the request must fail and the segment users cache must not be refreshed
889+
// (no ListSegmentUsersBySegment/Put expectations are registered, so the mock
890+
// controller fails the test if they are called).
891+
func TestUpdateSegmentPublishFailureMySQL(t *testing.T) {
892+
t.Parallel()
893+
mockController := gomock.NewController(t)
894+
defer mockController.Finish()
895+
896+
ctx, cancel := context.WithCancel(context.Background())
897+
defer cancel()
898+
ctx = metadata.NewIncomingContext(ctx, metadata.MD{
899+
"accept-language": []string{"ja"},
900+
})
901+
ctx = setToken(ctx)
902+
903+
service := createFeatureService(mockController)
904+
domainPublisher := publishermock.NewMockPublisher(mockController)
905+
service.domainPublisher = domainPublisher
906+
service.dbClient.(*databasemock.MockClient).EXPECT().RunInTransactionV2(
907+
gomock.Any(), gomock.Any(),
908+
).DoAndReturn(func(ctx context.Context, fn func(ctx context.Context) error) error {
909+
return fn(ctx)
910+
})
911+
service.segmentStorage.(*storagemock.MockSegmentStorage).EXPECT().GetSegment(
912+
gomock.Any(), gomock.Any(), gomock.Any(),
913+
).Return(&domain.Segment{
914+
Segment: &featureproto.Segment{
915+
Id: "id0",
916+
},
917+
}, nil, nil)
918+
service.segmentStorage.(*storagemock.MockSegmentStorage).EXPECT().UpdateSegment(
919+
gomock.Any(), gomock.Any(), gomock.Any(),
920+
).Return(nil)
921+
domainPublisher.EXPECT().Publish(gomock.Any(), gomock.Any()).Return(errors.New("publish failed"))
922+
923+
_, err := service.UpdateSegment(ctx, &featureproto.UpdateSegmentRequest{
924+
Id: "id0",
925+
EnvironmentId: "ns0",
926+
Rules: &featureproto.RuleListValue{
927+
Values: []*featureproto.Rule{
928+
{
929+
Clauses: []*featureproto.Clause{
930+
{
931+
Attribute: "plan",
932+
Operator: featureproto.Clause_EQUALS,
933+
Values: []string{"premium"},
934+
},
935+
},
936+
},
937+
},
938+
},
939+
})
940+
assert.Error(t, err)
941+
}
942+
735943
func setToken(ctx context.Context) context.Context {
736944
t := &token.AccessToken{
737945
Issuer: "issuer",

0 commit comments

Comments
 (0)