From 2b496042a5bfab15ce2d5c9a959ae54c53855578 Mon Sep 17 00:00:00 2001 From: vagarwal-viant Date: Fri, 7 Aug 2026 11:14:52 -0700 Subject: [PATCH 1/2] ENG-00000 fix data race on Statelet.Filters --- service/reader/service.go | 6 +++--- service/session/state.go | 4 +--- view/state.go | 16 +++++++++++++--- view/state_test.go | 37 +++++++++++++++++++++++++++++++++++++ 4 files changed, 54 insertions(+), 9 deletions(-) create mode 100644 view/state_test.go diff --git a/service/reader/service.go b/service/reader/service.go index 28c49b3e..468e7967 100644 --- a/service/reader/service.go +++ b/service/reader/service.go @@ -628,9 +628,9 @@ func (s *Service) warmupMatcher(ctx context.Context, aView *view.View, statelet } } } - cloned := *statelet + cloned := statelet.CloneForSummary() cloned.Template = clonedTemplate - ok, err := applyWarmupIdentityProjection(aView, &cloned) + ok, err := applyWarmupIdentityProjection(aView, cloned) if err != nil { return nil, err } @@ -638,7 +638,7 @@ func (s *Service) warmupMatcher(ctx context.Context, aView *view.View, statelet return nil, nil } - matcher, err := s.sqlBuilder.CacheSQLWithOptions(ctx, aView, &cloned, nil, nil, parent) + matcher, err := s.sqlBuilder.CacheSQLWithOptions(ctx, aView, cloned, nil, nil, parent) if err != nil || matcher == nil { return matcher, err } diff --git a/service/session/state.go b/service/session/state.go index c00b7428..91999460 100644 --- a/service/session/state.go +++ b/service/session/state.go @@ -73,9 +73,7 @@ func (s *Session) NewSession(component *repository.Component) *Session { if ret.Options.state != nil { ret.Options.state.RWMutex.Lock() for _, st := range ret.Options.state.Views { - if st != nil { - st.Filters = nil - } + st.ClearFilters() } ret.Options.state.RWMutex.Unlock() } diff --git a/view/state.go b/view/state.go index 7aaee716..6bc6b6e7 100644 --- a/view/state.go +++ b/view/state.go @@ -88,6 +88,16 @@ func (s *Statelet) AppendFilters(filters predicate.Filters) { s.filtersMu.Unlock() } +// ClearFilters safely clears the selector's filters. +func (s *Statelet) ClearFilters() { + if s == nil { + return + } + s.filtersMu.Lock() + s.Filters = nil + s.filtersMu.Unlock() +} + // NewStatelet creates a selector func NewStatelet() *Statelet { return &Statelet{ @@ -186,9 +196,9 @@ func (s *Statelet) CloneForSummary() *Statelet { ret._columnNames = map[string]bool{} } - if len(s.Filters) > 0 { - ret.Filters = append(predicate.Filters(nil), s.Filters...) - } + s.filtersMu.Lock() + ret.Filters = append(predicate.Filters(nil), s.Filters...) + s.filtersMu.Unlock() if len(s.Fields) > 0 { ret.Fields = append([]string(nil), s.Fields...) diff --git a/view/state_test.go b/view/state_test.go new file mode 100644 index 00000000..51814ee0 --- /dev/null +++ b/view/state_test.go @@ -0,0 +1,37 @@ +package view + +import ( + "sync" + "testing" + + "github.com/viant/datly/view/state/predicate" +) + +func TestStateletCloneForSummaryConcurrentFilters(t *testing.T) { + statelet := NewStatelet() + filter := &predicate.Filter{Name: "active"} + + var waitGroup sync.WaitGroup + waitGroup.Add(2) + + go func() { + defer waitGroup.Done() + for i := 0; i < 1000; i++ { + statelet.AppendFilters(predicate.Filters{filter}) + statelet.ClearFilters() + } + }() + + go func() { + defer waitGroup.Done() + for i := 0; i < 1000; i++ { + clone := statelet.CloneForSummary() + if clone == statelet { + t.Errorf("CloneForSummary() returned the original statelet") + return + } + } + }() + + waitGroup.Wait() +} From b9c5d49e6d2924bc2b868d362a253eb3a24b5b20 Mon Sep 17 00:00:00 2001 From: vagarwal-viant Date: Fri, 7 Aug 2026 11:27:57 -0700 Subject: [PATCH 2/2] ENG-00000 fix npe on readRequestBody(nil) --- view/state/kind/locator/body.go | 9 ++++-- view/state/kind/locator/body_test.go | 44 ++++++++++++++++++++++++++++ view/state/kind/locator/http.go | 3 ++ 3 files changed, 53 insertions(+), 3 deletions(-) create mode 100644 view/state/kind/locator/body_test.go diff --git a/view/state/kind/locator/body.go b/view/state/kind/locator/body.go index 6374362a..d98fe3e5 100644 --- a/view/state/kind/locator/body.go +++ b/view/state/kind/locator/body.go @@ -58,12 +58,12 @@ func (r *Body) Value(ctx context.Context, rType reflect.Type, name string) (inte } } - if len(r.body) == 0 { - return nil, false, nil - } if r.err != nil { return nil, false, r.err } + if len(r.body) == 0 { + return nil, false, nil + } if r.bodyType.Kind() == reflect.Map { return r.decodeBodyMap(ctx) } @@ -100,6 +100,9 @@ func (r *Body) initOnce() { // Non-multipart: clone and read body safely var request *http.Request request, r.err = shared.CloneHTTPRequest(r.request) + if r.err != nil { + return + } r.body, r.err = readRequestBody(request) }) } diff --git a/view/state/kind/locator/body_test.go b/view/state/kind/locator/body_test.go new file mode 100644 index 00000000..0f9e21e7 --- /dev/null +++ b/view/state/kind/locator/body_test.go @@ -0,0 +1,44 @@ +package locator + +import ( + "context" + "errors" + "net/http/httptest" + "reflect" + "testing" + + "github.com/stretchr/testify/require" +) + +type failingBody struct{} + +func (failingBody) Read([]byte) (int, error) { + return 0, errors.New("failed to read request body") +} + +func (failingBody) Close() error { + return nil +} + +func TestBodyValueReturnsRequestBodyReadError(t *testing.T) { + request := httptest.NewRequest("POST", "http://localhost/test", nil) + request.Body = failingBody{} + + aLocator, err := NewBody( + WithRequest(request), + WithBodyType(reflect.TypeOf(struct{}{})), + WithUnmarshal(func([]byte, interface{}) error { return nil }), + ) + require.NoError(t, err) + + value, ok, err := aLocator.Value(context.Background(), reflect.TypeOf(struct{}{}), "") + require.Nil(t, value) + require.False(t, ok) + require.EqualError(t, err, "failed to read request body") +} + +func TestReadRequestBodyRejectsNilRequest(t *testing.T) { + data, err := readRequestBody(nil) + require.Nil(t, data) + require.EqualError(t, err, "request was empty") +} diff --git a/view/state/kind/locator/http.go b/view/state/kind/locator/http.go index 62c9ed69..e76f9546 100644 --- a/view/state/kind/locator/http.go +++ b/view/state/kind/locator/http.go @@ -61,6 +61,9 @@ func NewHttpRequest(opts ...Option) (kind.Locator, error) { } func readRequestBody(request *http.Request) ([]byte, error) { + if request == nil { + return nil, fmt.Errorf("request was empty") + } if request.Body == nil { return nil, nil }