Skip to content

Commit 8054076

Browse files
fix: do not try to merge the already consolidated delete requests while listing them (#18544)
1 parent 5e28950 commit 8054076

5 files changed

Lines changed: 25 additions & 119 deletions

File tree

‎pkg/compactor/deletion/delete_requests_store_test.go‎

Lines changed: 23 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -173,8 +173,19 @@ func TestBatchCreateGetAllStoreTypes(t *testing.T) {
173173
tc := setupStoreType(t, storeType)
174174
defer tc.store.Stop()
175175

176-
reqID, err := tc.store.AddDeleteRequest(context.Background(), user1, `{foo="bar"}`, now.Add(-24*time.Hour), now, time.Hour)
176+
deletionQuery := `{foo="bar"}`
177+
reqID, err := tc.store.AddDeleteRequest(context.Background(), user1, deletionQuery, now.Add(-24*time.Hour), now, time.Hour)
178+
require.NoError(t, err)
179+
180+
// GetAllDeleteRequestsForUser should list consolidated requests and not shards
181+
consolidatedRequests, err := tc.store.GetAllDeleteRequestsForUser(context.Background(), user1, false)
177182
require.NoError(t, err)
183+
require.Len(t, consolidatedRequests, 1)
184+
require.Equal(t, consolidatedRequests[0].RequestID, reqID)
185+
require.Equal(t, consolidatedRequests[0].StartTime, now.Add(-24*time.Hour))
186+
require.Equal(t, consolidatedRequests[0].EndTime, now)
187+
require.Equal(t, consolidatedRequests[0].Query, deletionQuery)
188+
require.Equal(t, consolidatedRequests[0].Status, StatusReceived)
178189

179190
savedRequests, err := tc.store.GetUnprocessedShards(context.Background())
180191
require.NoError(t, err)
@@ -191,6 +202,17 @@ func TestBatchCreateGetAllStoreTypes(t *testing.T) {
191202
req, err := tc.store.GetDeleteRequest(context.Background(), user1, reqID)
192203
require.NoError(t, err)
193204
require.Equal(t, deleteRequestStatus(1, len(savedRequests)), req.Status)
205+
206+
// check that we get appropriate status in the consolidated request
207+
consolidatedRequests, err = tc.store.GetAllDeleteRequestsForUser(context.Background(), user1, false)
208+
require.NoError(t, err)
209+
require.Len(t, consolidatedRequests, 1)
210+
require.Equal(t, consolidatedRequests[0].RequestID, reqID)
211+
require.Equal(t, consolidatedRequests[0].StartTime, now.Add(-24*time.Hour))
212+
require.Equal(t, consolidatedRequests[0].EndTime, now)
213+
require.Equal(t, consolidatedRequests[0].Query, deletionQuery)
214+
require.Equal(t, consolidatedRequests[0].Status, deleteRequestStatus(1, len(savedRequests)))
215+
194216
})
195217

196218
t.Run("deletes several delete requests", func(t *testing.T) {

‎pkg/compactor/deletion/grpc_request_handler.go‎

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -39,14 +39,12 @@ func (g *GRPCRequestHandler) GetDeleteRequests(ctx context.Context, req *grpc.Ge
3939
return nil, errors.New(deletionNotAvailableMsg)
4040
}
4141

42-
deleteGroups, err := g.deleteRequestsStore.GetAllDeleteRequestsForUser(ctx, userID, req.ForQuerytimeFiltering)
42+
deleteRequests, err := g.deleteRequestsStore.GetAllDeleteRequestsForUser(ctx, userID, req.ForQuerytimeFiltering)
4343
if err != nil {
4444
level.Error(util_log.Logger).Log("msg", "error getting delete requests from the store", "err", err)
4545
return nil, err
4646
}
4747

48-
deleteRequests := mergeDeletes(deleteGroups)
49-
5048
sort.Slice(deleteRequests, func(i, j int) bool {
5149
return deleteRequests[i].CreatedAt < deleteRequests[j].CreatedAt
5250
})

‎pkg/compactor/deletion/grpc_request_handler_test.go‎

Lines changed: 0 additions & 55 deletions
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,6 @@ import (
55
"errors"
66
"net"
77
"testing"
8-
"time"
98

109
"github.com/grafana/dskit/middleware"
1110
"github.com/grafana/dskit/tenant"
@@ -89,60 +88,6 @@ func TestGRPCGetDeleteRequests(t *testing.T) {
8988
require.ElementsMatch(t, store.getAllResult, grpcDeleteRequestsToDeleteRequests(resp.DeleteRequests))
9089
})
9190

92-
t.Run("it merges requests with the same requestID", func(t *testing.T) {
93-
store := &mockDeleteRequestsStore{}
94-
store.getAllResult = []DeleteRequest{
95-
{RequestID: "test-request-1", CreatedAt: now, StartTime: now, EndTime: now.Add(time.Hour)},
96-
{RequestID: "test-request-1", CreatedAt: now, StartTime: now.Add(2 * time.Hour), EndTime: now.Add(3 * time.Hour)},
97-
{RequestID: "test-request-2", CreatedAt: now.Add(time.Minute), StartTime: now.Add(30 * time.Minute), EndTime: now.Add(90 * time.Minute)},
98-
{RequestID: "test-request-1", CreatedAt: now, StartTime: now.Add(time.Hour), EndTime: now.Add(2 * time.Hour)},
99-
}
100-
h := NewGRPCRequestHandler(store, &fakeLimits{defaultLimit: limit{deletionMode: deletionmode.FilterAndDelete.String()}})
101-
grpcClient, closer := server(t, h)
102-
t.Cleanup(closer)
103-
104-
ctx, _ := user.InjectIntoGRPCRequest(user.InjectOrgID(context.Background(), user1))
105-
orgID, err := tenant.TenantID(ctx)
106-
require.NoError(t, err)
107-
require.Equal(t, user1, orgID)
108-
109-
resp, err := grpcClient.GetDeleteRequests(ctx, &compactor_client_grpc.GetDeleteRequestsRequest{})
110-
require.NoError(t, err)
111-
require.ElementsMatch(t, []DeleteRequest{
112-
{RequestID: "test-request-1", Status: StatusReceived, CreatedAt: now, StartTime: now, EndTime: now.Add(3 * time.Hour)},
113-
{RequestID: "test-request-2", Status: StatusReceived, CreatedAt: now.Add(time.Minute), StartTime: now.Add(30 * time.Minute), EndTime: now.Add(90 * time.Minute)},
114-
}, grpcDeleteRequestsToDeleteRequests(resp.DeleteRequests))
115-
})
116-
117-
t.Run("it only considers a request processed if all it's subqueries are processed", func(t *testing.T) {
118-
store := &mockDeleteRequestsStore{}
119-
store.getAllResult = []DeleteRequest{
120-
{RequestID: "test-request-1", CreatedAt: now, Status: StatusProcessed},
121-
{RequestID: "test-request-1", CreatedAt: now, Status: StatusReceived},
122-
{RequestID: "test-request-1", CreatedAt: now, Status: StatusProcessed},
123-
{RequestID: "test-request-2", CreatedAt: now.Add(time.Minute), Status: StatusProcessed},
124-
{RequestID: "test-request-2", CreatedAt: now.Add(time.Minute), Status: StatusProcessed},
125-
{RequestID: "test-request-2", CreatedAt: now.Add(time.Minute), Status: StatusProcessed},
126-
{RequestID: "test-request-3", CreatedAt: now.Add(2 * time.Minute), Status: StatusReceived},
127-
}
128-
h := NewGRPCRequestHandler(store, &fakeLimits{defaultLimit: limit{deletionMode: deletionmode.FilterAndDelete.String()}})
129-
grpcClient, closer := server(t, h)
130-
t.Cleanup(closer)
131-
132-
ctx, _ := user.InjectIntoGRPCRequest(user.InjectOrgID(context.Background(), user1))
133-
orgID, err := tenant.TenantID(ctx)
134-
require.NoError(t, err)
135-
require.Equal(t, user1, orgID)
136-
137-
resp, err := grpcClient.GetDeleteRequests(ctx, &compactor_client_grpc.GetDeleteRequestsRequest{})
138-
require.NoError(t, err)
139-
require.ElementsMatch(t, []DeleteRequest{
140-
{RequestID: "test-request-1", CreatedAt: now, Status: "66% Complete"},
141-
{RequestID: "test-request-2", CreatedAt: now.Add(time.Minute), Status: StatusProcessed},
142-
{RequestID: "test-request-3", CreatedAt: now.Add(2 * time.Minute), Status: StatusReceived},
143-
}, grpcDeleteRequestsToDeleteRequests(resp.DeleteRequests))
144-
})
145-
14691
t.Run("error getting from store", func(t *testing.T) {
14792
store := &mockDeleteRequestsStore{}
14893
store.getAllErr = errors.New("something bad")

‎pkg/compactor/deletion/request_handler.go‎

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -143,14 +143,13 @@ func (dm *DeleteRequestHandler) GetAllDeleteRequestsHandler(w http.ResponseWrite
143143
}
144144

145145
forQuerytimeFiltering := r.URL.Query().Get(ForQuerytimeFilteringQueryParam) == "true"
146-
deleteGroups, err := dm.deleteRequestsStore.GetAllDeleteRequestsForUser(ctx, userID, forQuerytimeFiltering)
146+
deleteRequests, err := dm.deleteRequestsStore.GetAllDeleteRequestsForUser(ctx, userID, forQuerytimeFiltering)
147147
if err != nil {
148148
level.Error(util_log.Logger).Log("msg", "error getting delete requests from the store", "err", err)
149149
http.Error(w, err.Error(), http.StatusInternalServerError)
150150
return
151151
}
152152

153-
deleteRequests := mergeDeletes(deleteGroups)
154153
// We have to retain UserID and SequenceNum in json encoding of deletion manifests.
155154
// However, we do not want to return these to the users since they are not relevant.
156155
// ToDo(Sandeep): See if we can avoid doing it when we move to proto encoding for deletion manifests.

‎pkg/compactor/deletion/request_handler_test.go‎

Lines changed: 0 additions & 58 deletions
Original file line numberDiff line numberDiff line change
@@ -339,64 +339,6 @@ func TestGetAllDeleteRequestsHandler(t *testing.T) {
339339
}
340340
})
341341

342-
t.Run("it merges requests with the same requestID", func(t *testing.T) {
343-
store := &mockDeleteRequestsStore{}
344-
store.getAllResult = []DeleteRequest{
345-
{RequestID: "test-request-1", CreatedAt: now, StartTime: now, EndTime: now.Add(time.Hour)},
346-
{RequestID: "test-request-1", CreatedAt: now, StartTime: now.Add(2 * time.Hour), EndTime: now.Add(3 * time.Hour)},
347-
{RequestID: "test-request-2", CreatedAt: now.Add(time.Minute), StartTime: now.Add(30 * time.Minute), EndTime: now.Add(90 * time.Minute)},
348-
{RequestID: "test-request-1", CreatedAt: now, StartTime: now.Add(time.Hour), EndTime: now.Add(2 * time.Hour)},
349-
}
350-
h := NewDeleteRequestHandler(store, 0, 0, nil)
351-
352-
req := buildRequest("org-id", ``, "", "", false)
353-
354-
w := httptest.NewRecorder()
355-
h.GetAllDeleteRequestsHandler(w, req)
356-
357-
require.Equal(t, w.Code, http.StatusOK)
358-
359-
var result []DeleteRequest
360-
require.NoError(t, json.Unmarshal(w.Body.Bytes(), &result))
361-
362-
require.Len(t, result, 2)
363-
require.Equal(t, []DeleteRequest{
364-
{RequestID: "test-request-1", Status: StatusReceived, CreatedAt: now, StartTime: now, EndTime: now.Add(3 * time.Hour)},
365-
{RequestID: "test-request-2", Status: StatusReceived, CreatedAt: now.Add(time.Minute), StartTime: now.Add(30 * time.Minute), EndTime: now.Add(90 * time.Minute)},
366-
}, result)
367-
})
368-
369-
t.Run("it only considers a request processed if all it's subqueries are processed", func(t *testing.T) {
370-
store := &mockDeleteRequestsStore{}
371-
store.getAllResult = []DeleteRequest{
372-
{RequestID: "test-request-1", CreatedAt: now, Status: StatusProcessed},
373-
{RequestID: "test-request-1", CreatedAt: now, Status: StatusReceived},
374-
{RequestID: "test-request-1", CreatedAt: now, Status: StatusProcessed},
375-
{RequestID: "test-request-2", CreatedAt: now.Add(time.Minute), Status: StatusProcessed},
376-
{RequestID: "test-request-2", CreatedAt: now.Add(time.Minute), Status: StatusProcessed},
377-
{RequestID: "test-request-2", CreatedAt: now.Add(time.Minute), Status: StatusProcessed},
378-
{RequestID: "test-request-3", CreatedAt: now.Add(2 * time.Minute), Status: StatusReceived},
379-
}
380-
h := NewDeleteRequestHandler(store, 0, 0, nil)
381-
382-
req := buildRequest("org-id", ``, "", "", false)
383-
384-
w := httptest.NewRecorder()
385-
h.GetAllDeleteRequestsHandler(w, req)
386-
387-
require.Equal(t, w.Code, http.StatusOK)
388-
389-
var result []DeleteRequest
390-
require.NoError(t, json.Unmarshal(w.Body.Bytes(), &result))
391-
392-
require.Len(t, result, 3)
393-
require.Equal(t, []DeleteRequest{
394-
{RequestID: "test-request-1", CreatedAt: now, Status: "66% Complete"},
395-
{RequestID: "test-request-2", CreatedAt: now.Add(time.Minute), Status: StatusProcessed},
396-
{RequestID: "test-request-3", CreatedAt: now.Add(2 * time.Minute), Status: StatusReceived},
397-
}, result)
398-
})
399-
400342
t.Run("error getting from store", func(t *testing.T) {
401343
store := &mockDeleteRequestsStore{}
402344
store.getAllErr = errors.New("something bad")

0 commit comments

Comments
 (0)