diff --git a/cmd/e2e/api_test.go b/cmd/e2e/api_test.go index 27384e1f6..6b02c54f4 100644 --- a/cmd/e2e/api_test.go +++ b/cmd/e2e/api_test.go @@ -534,9 +534,9 @@ func (suite *basicSuite) TestListTenantsAPI() { body, ok := resp.Body.(map[string]interface{}) require.True(t, ok, "response should be a map") - data, ok := body["data"].([]interface{}) - require.True(t, ok, "data should be an array") - assert.GreaterOrEqual(t, len(data), 3, "should have at least 3 tenants") + models, ok := body["models"].([]interface{}) + require.True(t, ok, "models should be an array") + assert.GreaterOrEqual(t, len(models), 3, "should have at least 3 tenants") }) // Test list with limit @@ -550,9 +550,9 @@ func (suite *basicSuite) TestListTenantsAPI() { body, ok := resp.Body.(map[string]interface{}) require.True(t, ok, "response should be a map") - data, ok := body["data"].([]interface{}) - require.True(t, ok, "data should be an array") - assert.Equal(t, 2, len(data), "should have exactly 2 tenants") + models, ok := body["models"].([]interface{}) + require.True(t, ok, "models should be an array") + assert.Equal(t, 2, len(models), "should have exactly 2 tenants") }) // Test invalid limit @@ -577,11 +577,12 @@ func (suite *basicSuite) TestListTenantsAPI() { body, ok := resp.Body.(map[string]interface{}) require.True(t, ok, "response should be a map") - data, ok := body["data"].([]interface{}) - require.True(t, ok, "data should be an array") - assert.Equal(t, 2, len(data), "page 1 should have 2 tenants") + models, ok := body["models"].([]interface{}) + require.True(t, ok, "models should be an array") + assert.Equal(t, 2, len(models), "page 1 should have 2 tenants") - next, _ := body["next"].(string) + pagination, _ := body["pagination"].(map[string]interface{}) + next, _ := pagination["next"].(string) require.NotEmpty(t, next, "should have next cursor") // Get second page using next cursor @@ -594,11 +595,12 @@ func (suite *basicSuite) TestListTenantsAPI() { body, ok = resp.Body.(map[string]interface{}) require.True(t, ok, "response should be a map") - data, ok = body["data"].([]interface{}) - require.True(t, ok, "data should be an array") - assert.GreaterOrEqual(t, len(data), 1, "page 2 should have at least 1 tenant") + models, ok = body["models"].([]interface{}) + require.True(t, ok, "models should be an array") + assert.GreaterOrEqual(t, len(models), 1, "page 2 should have at least 1 tenant") - prev, _ := body["prev"].(string) + pagination, _ = body["pagination"].(map[string]interface{}) + prev, _ := pagination["prev"].(string) assert.NotEmpty(t, prev, "page 2 should have prev cursor") }) @@ -613,7 +615,8 @@ func (suite *basicSuite) TestListTenantsAPI() { body, ok := resp.Body.(map[string]interface{}) require.True(t, ok) - next, _ := body["next"].(string) + pagination, _ := body["pagination"].(map[string]interface{}) + next, _ := pagination["next"].(string) require.NotEmpty(t, next, "should have next cursor") // Go to page 2 @@ -625,7 +628,8 @@ func (suite *basicSuite) TestListTenantsAPI() { body, ok = resp.Body.(map[string]interface{}) require.True(t, ok) - prev, _ := body["prev"].(string) + pagination, _ = body["pagination"].(map[string]interface{}) + prev, _ := pagination["prev"].(string) require.NotEmpty(t, prev, "page 2 should have prev cursor") // Using prev cursor returns items with newer timestamps (keyset pagination) @@ -639,9 +643,9 @@ func (suite *basicSuite) TestListTenantsAPI() { body, ok = resp.Body.(map[string]interface{}) require.True(t, ok, "response should be a map") - data, ok := body["data"].([]interface{}) - require.True(t, ok, "data should be an array") - assert.NotEmpty(t, data, "prev cursor should return items") + models, ok := body["models"].([]interface{}) + require.True(t, ok, "models should be an array") + assert.NotEmpty(t, models, "prev cursor should return items") }) // Cleanup @@ -727,9 +731,9 @@ func (suite *basicSuite) TestDestinationsAPI() { StatusCode: http.StatusUnprocessableEntity, Body: map[string]interface{}{ "message": "validation error", - "data": map[string]interface{}{ - "topics": "required", - "type": "required", + "data": []interface{}{ + "type is required", + "topics is required", }, }, }, @@ -795,8 +799,8 @@ func (suite *basicSuite) TestDestinationsAPI() { StatusCode: http.StatusUnprocessableEntity, Body: map[string]interface{}{ "message": "validation error", - "data": map[string]interface{}{ - "config.url": "required", + "data": []interface{}{ + "config.url is required", }, }, }, @@ -1253,8 +1257,8 @@ func (suite *basicSuite) TestDestinationsAPI() { StatusCode: http.StatusUnprocessableEntity, Body: map[string]interface{}{ "message": "validation error", - "data": map[string]interface{}{ - "config.url": "required", + "data": []interface{}{ + "config.url is required", }, }, }, diff --git a/cmd/e2e/destwebhook_test.go b/cmd/e2e/destwebhook_test.go index 4aa26191f..7d224f00a 100644 --- a/cmd/e2e/destwebhook_test.go +++ b/cmd/e2e/destwebhook_test.go @@ -530,8 +530,8 @@ func (suite *basicSuite) TestDestwebhookTenantSecretManagement() { StatusCode: http.StatusUnprocessableEntity, Body: map[string]interface{}{ "message": "validation error", - "data": map[string]interface{}{ - "credentials.secret": "forbidden", + "data": []interface{}{ + "credentials.secret failed forbidden validation", }, }, }, @@ -593,8 +593,8 @@ func (suite *basicSuite) TestDestwebhookTenantSecretManagement() { StatusCode: http.StatusUnprocessableEntity, Body: map[string]interface{}{ "message": "validation error", - "data": map[string]interface{}{ - "credentials.secret": "forbidden", + "data": []interface{}{ + "credentials.secret failed forbidden validation", }, }, }, @@ -616,8 +616,8 @@ func (suite *basicSuite) TestDestwebhookTenantSecretManagement() { StatusCode: http.StatusUnprocessableEntity, Body: map[string]interface{}{ "message": "validation error", - "data": map[string]interface{}{ - "credentials.previous_secret": "forbidden", + "data": []interface{}{ + "credentials.previous_secret failed forbidden validation", }, }, }, @@ -639,8 +639,8 @@ func (suite *basicSuite) TestDestwebhookTenantSecretManagement() { StatusCode: http.StatusUnprocessableEntity, Body: map[string]interface{}{ "message": "validation error", - "data": map[string]interface{}{ - "credentials.previous_secret_invalid_at": "forbidden", + "data": []interface{}{ + "credentials.previous_secret_invalid_at failed forbidden validation", }, }, }, @@ -865,8 +865,8 @@ func (suite *basicSuite) TestDestwebhookAdminSecretManagement() { StatusCode: http.StatusUnprocessableEntity, Body: map[string]interface{}{ "message": "validation error", - "data": map[string]interface{}{ - "credentials.rotate_secret": "invalid", + "data": []interface{}{ + "credentials.rotate_secret failed invalid validation", }, }, }, @@ -932,8 +932,8 @@ func (suite *basicSuite) TestDestwebhookAdminSecretManagement() { StatusCode: http.StatusUnprocessableEntity, Body: map[string]interface{}{ "message": "validation error", - "data": map[string]interface{}{ - "credentials.previous_secret_invalid_at": "pattern", + "data": []interface{}{ + "credentials.previous_secret_invalid_at failed pattern validation", }, }, }, @@ -1082,8 +1082,8 @@ func (suite *basicSuite) TestDestwebhookAdminSecretManagement() { StatusCode: http.StatusUnprocessableEntity, Body: map[string]interface{}{ "message": "validation error", - "data": map[string]interface{}{ - "credentials.secret": "required", + "data": []interface{}{ + "credentials.secret is required", }, }, }, diff --git a/cmd/e2e/log_test.go b/cmd/e2e/log_test.go index c48e15d51..c6c1729aa 100644 --- a/cmd/e2e/log_test.go +++ b/cmd/e2e/log_test.go @@ -3,7 +3,6 @@ package e2e_test import ( "fmt" "net/http" - "net/url" "time" "github.com/hookdeck/outpost/cmd/e2e/httpclient" @@ -92,8 +91,10 @@ func (suite *basicSuite) TestLogAPI() { } suite.RunAPITests(suite.T(), setupTests) - // Publish 10 events with small delays for distinct timestamps + // Publish 10 events with explicit timestamps (1 second apart) + baseTime := time.Now().Add(-1 * time.Hour).Truncate(time.Second) for i, eventID := range eventIDs { + eventTime := baseTime.Add(time.Duration(i) * time.Second) resp, err := suite.client.Do(suite.AuthRequest(httpclient.Request{ Method: httpclient.MethodPOST, Path: "/publish", @@ -102,12 +103,12 @@ func (suite *basicSuite) TestLogAPI() { "tenant_id": tenantID, "topic": "user.created", "eligible_for_retry": true, + "time": eventTime.Format(time.RFC3339Nano), "data": map[string]interface{}{"index": i}, }, })) suite.Require().NoError(err) suite.Require().Equal(http.StatusAccepted, resp.StatusCode, "failed to publish event %d", i) - time.Sleep(50 * time.Millisecond) } // Wait for all deliveries (30s timeout for slow CI environments) @@ -126,11 +127,11 @@ func (suite *basicSuite) TestLogAPI() { suite.Require().Equal(http.StatusOK, resp.StatusCode) body := resp.Body.(map[string]interface{}) - data := body["data"].([]interface{}) - suite.Len(data, 10) + models := body["models"].([]interface{}) + suite.Len(models, 10) // Verify structure - first := data[0].(map[string]interface{}) + first := models[0].(map[string]interface{}) suite.NotEmpty(first["id"]) suite.NotEmpty(first["event"]) suite.Equal(destinationID, first["destination"]) @@ -147,8 +148,8 @@ func (suite *basicSuite) TestLogAPI() { suite.Require().Equal(http.StatusOK, resp.StatusCode) body := resp.Body.(map[string]interface{}) - data := body["data"].([]interface{}) - suite.Len(data, 10) + models := body["models"].([]interface{}) + suite.Len(models, 10) }) suite.Run("filter by event_id", func() { @@ -160,8 +161,8 @@ func (suite *basicSuite) TestLogAPI() { suite.Require().Equal(http.StatusOK, resp.StatusCode) body := resp.Body.(map[string]interface{}) - data := body["data"].([]interface{}) - suite.Len(data, 1) + models := body["models"].([]interface{}) + suite.Len(models, 1) }) suite.Run("include=event returns event object without data", func() { @@ -173,10 +174,10 @@ func (suite *basicSuite) TestLogAPI() { suite.Require().Equal(http.StatusOK, resp.StatusCode) body := resp.Body.(map[string]interface{}) - data := body["data"].([]interface{}) - suite.Require().Len(data, 1) + models := body["models"].([]interface{}) + suite.Require().Len(models, 1) - delivery := data[0].(map[string]interface{}) + delivery := models[0].(map[string]interface{}) event := delivery["event"].(map[string]interface{}) suite.NotEmpty(event["id"]) suite.NotEmpty(event["topic"]) @@ -193,10 +194,10 @@ func (suite *basicSuite) TestLogAPI() { suite.Require().Equal(http.StatusOK, resp.StatusCode) body := resp.Body.(map[string]interface{}) - data := body["data"].([]interface{}) - suite.Require().Len(data, 1) + models := body["models"].([]interface{}) + suite.Require().Len(models, 1) - delivery := data[0].(map[string]interface{}) + delivery := models[0].(map[string]interface{}) event := delivery["event"].(map[string]interface{}) suite.NotEmpty(event["id"]) suite.NotNil(event["data"]) // include=event.data SHOULD include data @@ -211,10 +212,10 @@ func (suite *basicSuite) TestLogAPI() { suite.Require().Equal(http.StatusOK, resp.StatusCode) body := resp.Body.(map[string]interface{}) - data := body["data"].([]interface{}) - suite.Require().Len(data, 1) + models := body["models"].([]interface{}) + suite.Require().Len(models, 1) - delivery := data[0].(map[string]interface{}) + delivery := models[0].(map[string]interface{}) suite.NotNil(delivery["response_data"]) }) }) @@ -232,11 +233,11 @@ func (suite *basicSuite) TestLogAPI() { suite.Require().Equal(http.StatusOK, resp.StatusCode) body := resp.Body.(map[string]interface{}) - data := body["data"].([]interface{}) - suite.Len(data, 10) + models := body["models"].([]interface{}) + suite.Len(models, 10) // Verify structure - first := data[0].(map[string]interface{}) + first := models[0].(map[string]interface{}) suite.NotEmpty(first["id"]) suite.NotEmpty(first["topic"]) suite.NotEmpty(first["time"]) @@ -252,8 +253,8 @@ func (suite *basicSuite) TestLogAPI() { suite.Require().Equal(http.StatusOK, resp.StatusCode) body := resp.Body.(map[string]interface{}) - data := body["data"].([]interface{}) - suite.Len(data, 10) // All events have topic=user.created + models := body["models"].([]interface{}) + suite.Len(models, 10) // All events have topic=user.created }) suite.Run("retrieve single event", func() { @@ -279,18 +280,18 @@ func (suite *basicSuite) TestLogAPI() { suite.Equal(http.StatusNotFound, resp.StatusCode) }) - suite.Run("filter by start time excludes past events", func() { - futureTime := url.QueryEscape(time.Now().Add(1 * time.Hour).Format(time.RFC3339)) + suite.Run("filter by time[gte] excludes past events", func() { + futureTime := time.Now().Add(1 * time.Hour).UTC().Format(time.RFC3339) resp, err := suite.client.Do(suite.AuthRequest(httpclient.Request{ Method: httpclient.MethodGET, - Path: "/tenants/" + tenantID + "/events?start=" + futureTime, + Path: "/tenants/" + tenantID + "/events?time[gte]=" + futureTime, })) suite.Require().NoError(err) suite.Require().Equal(http.StatusOK, resp.StatusCode) body := resp.Body.(map[string]interface{}) - data := body["data"].([]interface{}) - suite.Len(data, 0) + models := body["models"].([]interface{}) + suite.Len(models, 0) }) }) @@ -301,18 +302,18 @@ func (suite *basicSuite) TestLogAPI() { suite.Run("events desc returns newest first", func() { resp, err := suite.client.Do(suite.AuthRequest(httpclient.Request{ Method: httpclient.MethodGET, - Path: "/tenants/" + tenantID + "/events?sort_order=desc", + Path: "/tenants/" + tenantID + "/events?dir=desc", })) suite.Require().NoError(err) suite.Require().Equal(http.StatusOK, resp.StatusCode) body := resp.Body.(map[string]interface{}) - data := body["data"].([]interface{}) - suite.Require().Len(data, 10) + models := body["models"].([]interface{}) + suite.Require().Len(models, 10) - for i := 0; i < len(data)-1; i++ { - curr := parseTime(data[i].(map[string]interface{})["time"].(string)) - next := parseTime(data[i+1].(map[string]interface{})["time"].(string)) + for i := 0; i < len(models)-1; i++ { + curr := parseTime(models[i].(map[string]interface{})["time"].(string)) + next := parseTime(models[i+1].(map[string]interface{})["time"].(string)) suite.True(curr.After(next) || curr.Equal(next), "events not in descending order at index %d", i) } }) @@ -320,26 +321,26 @@ func (suite *basicSuite) TestLogAPI() { suite.Run("events asc returns oldest first", func() { resp, err := suite.client.Do(suite.AuthRequest(httpclient.Request{ Method: httpclient.MethodGET, - Path: "/tenants/" + tenantID + "/events?sort_order=asc", + Path: "/tenants/" + tenantID + "/events?dir=asc", })) suite.Require().NoError(err) suite.Require().Equal(http.StatusOK, resp.StatusCode) body := resp.Body.(map[string]interface{}) - data := body["data"].([]interface{}) - suite.Require().Len(data, 10) + models := body["models"].([]interface{}) + suite.Require().Len(models, 10) - for i := 0; i < len(data)-1; i++ { - curr := parseTime(data[i].(map[string]interface{})["time"].(string)) - next := parseTime(data[i+1].(map[string]interface{})["time"].(string)) + for i := 0; i < len(models)-1; i++ { + curr := parseTime(models[i].(map[string]interface{})["time"].(string)) + next := parseTime(models[i+1].(map[string]interface{})["time"].(string)) suite.True(curr.Before(next) || curr.Equal(next), "events not in ascending order at index %d", i) } }) - suite.Run("events invalid sort_order returns 422", func() { + suite.Run("events invalid dir returns 422", func() { resp, err := suite.client.Do(suite.AuthRequest(httpclient.Request{ Method: httpclient.MethodGET, - Path: "/tenants/" + tenantID + "/events?sort_order=invalid", + Path: "/tenants/" + tenantID + "/events?dir=invalid", })) suite.Require().NoError(err) suite.Equal(http.StatusUnprocessableEntity, resp.StatusCode) @@ -357,7 +358,7 @@ func (suite *basicSuite) TestLogAPI() { pageCount := 0 for { - path := "/tenants/" + tenantID + "/events?limit=3&sort_order=asc" + path := "/tenants/" + tenantID + "/events?limit=3&dir=asc" if nextCursor != "" { path += "&next=" + nextCursor } @@ -370,15 +371,16 @@ func (suite *basicSuite) TestLogAPI() { suite.Require().Equal(http.StatusOK, resp.StatusCode) body := resp.Body.(map[string]interface{}) - data := body["data"].([]interface{}) + models := body["models"].([]interface{}) pageCount++ - for _, item := range data { + for _, item := range models { event := item.(map[string]interface{}) allEventIDs = append(allEventIDs, event["id"].(string)) } - if next, ok := body["next"].(string); ok && next != "" { + pagination, _ := body["pagination"].(map[string]interface{}) + if next, ok := pagination["next"].(string); ok && next != "" { nextCursor = next } else { break @@ -393,6 +395,83 @@ func (suite *basicSuite) TestLogAPI() { suite.Equal(4, pageCount, "expected 4 pages (3+3+3+1)") suite.Len(allEventIDs, 10, "should have all 10 events") }) + + suite.Run("cursor pagination with time filter", func() { + // Get all events to establish a time window + resp, err := suite.client.Do(suite.AuthRequest(httpclient.Request{ + Method: httpclient.MethodGET, + Path: "/tenants/" + tenantID + "/events?dir=asc&limit=10", + })) + suite.Require().NoError(err) + suite.Require().Equal(http.StatusOK, resp.StatusCode) + + body := resp.Body.(map[string]interface{}) + models := body["models"].([]interface{}) + suite.Require().Len(models, 10) + + // Use the 3rd and 7th events to create a time window + event3 := models[2].(map[string]interface{}) + event7 := models[6].(map[string]interface{}) + timeGTE := event3["time"].(string) + timeLTE := event7["time"].(string) + timeGTEParsed := parseTime(timeGTE) + timeLTEParsed := parseTime(timeLTE) + + // Paginate within the time window with limit=2 + var windowEvents []map[string]interface{} + nextCursor := "" + pageCount := 0 + + for { + path := "/tenants/" + tenantID + "/events?dir=asc&limit=2" + path += "&time[gte]=" + timeGTE + "&time[lte]=" + timeLTE + if nextCursor != "" { + path += "&next=" + nextCursor + } + + resp, err := suite.client.Do(suite.AuthRequest(httpclient.Request{ + Method: httpclient.MethodGET, + Path: path, + })) + suite.Require().NoError(err) + suite.Require().Equal(http.StatusOK, resp.StatusCode) + + body := resp.Body.(map[string]interface{}) + windowModels := body["models"].([]interface{}) + pageCount++ + + for _, item := range windowModels { + event := item.(map[string]interface{}) + windowEvents = append(windowEvents, event) + } + + pagination, _ := body["pagination"].(map[string]interface{}) + if next, ok := pagination["next"].(string); ok && next != "" { + nextCursor = next + } else { + break + } + + if pageCount > 10 { + suite.Fail("too many pages") + break + } + } + + // Verify time filter worked: should have fewer events than total + suite.Greater(len(windowEvents), 0, "should have some events in window") + suite.Less(len(windowEvents), 10, "time filter should exclude some events") + + // Verify pagination worked: multiple pages needed + suite.Greater(pageCount, 1, "should require multiple pages") + + // Verify all returned events are within the time window + for _, event := range windowEvents { + eventTime := parseTime(event["time"].(string)) + suite.True(!eventTime.Before(timeGTEParsed), "event time %v should be >= %v", eventTime, timeGTEParsed) + suite.True(!eventTime.After(timeLTEParsed), "event time %v should be <= %v", eventTime, timeLTEParsed) + } + }) }) // Cleanup @@ -535,9 +614,9 @@ func (suite *basicSuite) TestRetryAPI() { suite.Require().Equal(http.StatusOK, deliveriesResp.StatusCode) body := deliveriesResp.Body.(map[string]interface{}) - data := body["data"].([]interface{}) - suite.Require().NotEmpty(data, "should have at least one delivery") - firstDelivery := data[0].(map[string]interface{}) + models := body["models"].([]interface{}) + suite.Require().NotEmpty(models, "should have at least one delivery") + firstDelivery := models[0].(map[string]interface{}) deliveryID := firstDelivery["id"].(string) // Update mock to succeed for retry @@ -621,7 +700,7 @@ func (suite *basicSuite) TestRetryAPI() { "body": map[string]interface{}{ "type": "object", "properties": map[string]interface{}{ - "data": map[string]interface{}{ + "models": map[string]interface{}{ "type": "array", "minItems": 2, // Original + retry }, @@ -1208,13 +1287,13 @@ func (suite *basicSuite) TestAdminLogEndpoints() { suite.Require().Equal(http.StatusOK, resp.StatusCode) body := resp.Body.(map[string]interface{}) - data := body["data"].([]interface{}) + models := body["models"].([]interface{}) // Should have at least 2 events (one from each tenant we created) - suite.GreaterOrEqual(len(data), 2) + suite.GreaterOrEqual(len(models), 2) // Verify we have events from both tenants by checking event IDs eventsSeen := map[string]bool{} - for _, item := range data { + for _, item := range models { event := item.(map[string]interface{}) if id, ok := event["id"].(string); ok { eventsSeen[id] = true @@ -1233,13 +1312,13 @@ func (suite *basicSuite) TestAdminLogEndpoints() { suite.Require().Equal(http.StatusOK, resp.StatusCode) body := resp.Body.(map[string]interface{}) - data := body["data"].([]interface{}) + models := body["models"].([]interface{}) // Should have at least 2 deliveries (one from each tenant we created) - suite.GreaterOrEqual(len(data), 2) + suite.GreaterOrEqual(len(models), 2) // Verify we have deliveries from both tenants by checking event IDs eventsSeen := map[string]bool{} - for _, item := range data { + for _, item := range models { delivery := item.(map[string]interface{}) if event, ok := delivery["event"].(map[string]interface{}); ok { if id, ok := event["id"].(string); ok { @@ -1265,11 +1344,11 @@ func (suite *basicSuite) TestAdminLogEndpoints() { suite.Require().Equal(http.StatusOK, resp.StatusCode) body := resp.Body.(map[string]interface{}) - data := body["data"].([]interface{}) - suite.Len(data, 1) + models := body["models"].([]interface{}) + suite.Len(models, 1) // Verify only tenant1 event by ID - event := data[0].(map[string]interface{}) + event := models[0].(map[string]interface{}) suite.Equal(event1ID, event["id"]) }) @@ -1282,11 +1361,11 @@ func (suite *basicSuite) TestAdminLogEndpoints() { suite.Require().Equal(http.StatusOK, resp.StatusCode) body := resp.Body.(map[string]interface{}) - data := body["data"].([]interface{}) - suite.Len(data, 1) + models := body["models"].([]interface{}) + suite.Len(models, 1) // Verify only tenant2 delivery by event ID - delivery := data[0].(map[string]interface{}) + delivery := models[0].(map[string]interface{}) event := delivery["event"].(map[string]interface{}) suite.Equal(event2ID, event["id"]) }) diff --git a/cmd/e2e/suites_test.go b/cmd/e2e/suites_test.go index 37bd823d5..baf5bc3a8 100644 --- a/cmd/e2e/suites_test.go +++ b/cmd/e2e/suites_test.go @@ -61,8 +61,8 @@ func (s *e2eSuite) waitForDeliveries(t *testing.T, path string, minCount int, ti lastStatus = resp.StatusCode if resp.StatusCode == http.StatusOK { if body, ok := resp.Body.(map[string]interface{}); ok { - if data, ok := body["data"].([]interface{}); ok { - lastCount = len(data) + if models, ok := body["models"].([]interface{}); ok { + lastCount = len(models) if lastCount >= minCount { return } diff --git a/docs/apis/openapi.yaml b/docs/apis/openapi.yaml index 3d4bbb37d..39760b654 100644 --- a/docs/apis/openapi.yaml +++ b/docs/apis/openapi.yaml @@ -86,59 +86,17 @@ components: type: string nullable: true description: Optional metadata to store with the tenant. - TenantListItem: - type: object - description: Tenant object returned in list operations. - properties: - id: - type: string - description: User-defined system ID for the tenant. - example: "123" - destinations_count: - type: integer - description: Number of destinations associated with the tenant. - example: 5 - topics: - type: array - items: - type: string - description: List of subscribed topics across all destinations for this tenant. - example: ["user.created", "user.deleted"] - metadata: - type: object - additionalProperties: - type: string - nullable: true - description: Arbitrary key-value pairs for storing contextual information about the tenant. - created_at: - type: string - format: date-time - description: ISO Date when the tenant was created. - example: "2024-01-01T00:00:00Z" - updated_at: - type: string - format: date-time - description: ISO Date when the tenant was last updated. - example: "2024-01-01T00:00:00Z" - TenantListResponse: + TenantPaginatedResult: type: object description: Paginated list of tenants. properties: - data: + models: type: array items: - $ref: "#/components/schemas/TenantListItem" + $ref: "#/components/schemas/Tenant" description: Array of tenant objects. - next: - type: string - nullable: true - description: Cursor for the next page of results. Null if no more results. - example: "MTcwNDA2NzIwMA==" - prev: - type: string - nullable: true - description: Cursor for the previous page of results. Null if on first page. - example: null + pagination: + $ref: "#/components/schemas/SeekPagination" count: type: integer description: Total count of all tenants. @@ -198,26 +156,57 @@ components: customer: tier: "premium" - PaginatedResponse: + SeekPagination: type: object - required: [count, data, next, prev] + description: Cursor-based pagination metadata for list responses. properties: - count: + order_by: + type: string + description: The field being sorted on. + example: "created_at" + dir: + type: string + enum: [asc, desc] + description: Sort direction. + example: "desc" + limit: type: integer - description: Total number of items across all pages - example: 42 - data: - type: array - items: {} # Will be overridden by specific endpoints - description: Array of items for current page + description: Page size limit. + example: 100 next: type: string - description: Cursor for next page (empty string if no next page) - example: "" + nullable: true + description: Cursor for the next page of results. Null if no more results. + example: "MTcwNDA2NzIwMA==" prev: type: string - description: Cursor for previous page (empty string if no previous page) - example: "" + nullable: true + description: Cursor for the previous page of results. Null if on first page. + example: null + + APIErrorResponse: + type: object + description: Standard error response format. + properties: + status: + type: integer + description: HTTP status code. + example: 422 + message: + type: string + description: Human-readable error message. + example: "validation error" + data: + description: Additional error details. For validation errors, this is an array of human-readable messages. + oneOf: + - type: array + items: + type: string + description: Array of validation error messages. + example: ["email is required", "password must be at least 6 characters"] + - type: object + additionalProperties: true + description: Additional contextual data about the error. # Destination Type Specific Config/Credentials Schemas WebhookConfig: @@ -1688,6 +1677,11 @@ components: eligible_for_retry: type: boolean description: Should event delivery be retried on failure. + time: + type: string + format: date-time + description: Optional. Custom timestamp for the event. If not provided, defaults to the current time. + example: "2024-01-15T10:30:00Z" metadata: type: object description: Any key-value string pairs for metadata. @@ -1878,23 +1872,29 @@ components: additionalProperties: true description: The event payload data. example: { "user_id": "userid", "status": "active" } - DeliveryListResponse: + DeliveryPaginatedResult: type: object description: Paginated list of deliveries. properties: - data: + models: type: array items: $ref: "#/components/schemas/Delivery" description: Array of delivery objects. - next: - type: string - description: Cursor for the next page of results. Empty string if no more results. - example: "MTcwNDA2NzIwMA==" - prev: - type: string - description: Cursor for the previous page of results. Empty string if on first page. - example: "" + pagination: + $ref: "#/components/schemas/SeekPagination" + + EventPaginatedResult: + type: object + description: Paginated list of events. + properties: + models: + type: array + items: + $ref: "#/components/schemas/Event" + description: Array of event objects. + pagination: + $ref: "#/components/schemas/SeekPagination" # Destination Type Schema (for Metadata endpoint) DestinationType: @@ -2174,14 +2174,36 @@ paths: maximum: 100 default: 20 description: Number of tenants to return per page (1-100, default 20). - - name: order + - name: order_by + in: query + required: false + schema: + type: string + enum: [created_at] + default: created_at + description: Field to sort by. + - name: dir in: query required: false schema: type: string enum: [asc, desc] default: desc - description: Sort order by `created_at` timestamp. + description: Sort direction. + - name: created_at[gte] + in: query + required: false + schema: + type: string + format: date-time + description: Filter tenants created at or after this time (RFC3339 or YYYY-MM-DD format). + - name: created_at[lte] + in: query + required: false + schema: + type: string + format: date-time + description: Filter tenants created at or before this time (RFC3339 or YYYY-MM-DD format). - name: next in: query required: false @@ -2200,9 +2222,9 @@ paths: content: application/json: schema: - $ref: "#/components/schemas/TenantListResponse" + $ref: "#/components/schemas/TenantPaginatedResult" example: - data: + models: - id: "tenant_123" metadata: plan: "pro" @@ -2212,8 +2234,12 @@ paths: metadata: null created_at: "2024-01-14T09:00:00Z" updated_at: "2024-01-14T09:00:00Z" - next: "MTcwNDA2NzIwMA==" - prev: null + pagination: + order_by: "created_at" + dir: "desc" + limit: 20 + next: "MTcwNDA2NzIwMA==" + prev: null count: 42 "400": description: Invalid request parameters (e.g., invalid cursor, both next and prev provided). @@ -2341,20 +2367,20 @@ paths: items: type: string description: Filter events by topic(s). Can be specified multiple times or comma-separated. - - name: start + - name: time[gte] in: query required: false schema: type: string format: date-time - description: Filter events with time >= start (RFC3339 format). - - name: end + description: Filter events with time >= value (RFC3339 or YYYY-MM-DD format). + - name: time[lte] in: query required: false schema: type: string format: date-time - description: Filter events with time <= end (RFC3339 format). + description: Filter events with time <= value (RFC3339 or YYYY-MM-DD format). - name: limit in: query required: false @@ -2376,32 +2402,33 @@ paths: schema: type: string description: Cursor for previous page of results. - - name: sort_order + - name: order_by + in: query + required: false + schema: + type: string + enum: [time] + default: time + description: Field to sort by. + - name: dir in: query required: false schema: type: string enum: [asc, desc] default: desc - description: Sort order (ascending or descending). + description: Sort direction. responses: "200": description: A paginated list of events. content: application/json: schema: - allOf: - - $ref: "#/components/schemas/PaginatedResponse" - - type: object - properties: - data: - type: array - items: - $ref: "#/components/schemas/Event" + $ref: "#/components/schemas/EventPaginatedResult" examples: AdminEventsListExample: value: - data: + models: - id: "evt_123" topic: "user.created" time: "2024-01-01T00:00:00Z" @@ -2414,12 +2441,20 @@ paths: eligible_for_retry: true metadata: { "source": "oms" } data: { "order_id": "orderid", "tracking": "1Z..." } - next: "MTcwNDA2NzIwMA==" - prev: "" + pagination: + order_by: "time" + dir: "desc" + limit: 100 + next: "MTcwNDA2NzIwMA==" + prev: null "401": description: Unauthorized (Admin API Key missing or invalid). "422": description: Validation error (invalid query parameters). + content: + application/json: + schema: + $ref: "#/components/schemas/APIErrorResponse" /deliveries: get: @@ -2468,34 +2503,20 @@ paths: items: type: string description: Filter deliveries by event topic(s). Can be specified multiple times or comma-separated. - - name: start + - name: time[gte] in: query required: false schema: type: string format: date-time - description: Filter deliveries with delivered_at >= start (RFC3339 format). - - name: end + description: Filter deliveries by event time >= value (RFC3339 or YYYY-MM-DD format). + - name: time[lte] in: query required: false schema: type: string format: date-time - description: Filter deliveries with delivered_at <= end (RFC3339 format). - - name: event_start - in: query - required: false - schema: - type: string - format: date-time - description: Filter deliveries by event time >= event_start (RFC3339 format). - - name: event_end - in: query - required: false - schema: - type: string - format: date-time - description: Filter deliveries by event time <= event_end (RFC3339 format). + description: Filter deliveries by event time <= value (RFC3339 or YYYY-MM-DD format). - name: limit in: query required: false @@ -2531,33 +2552,33 @@ paths: - `event`: Include event summary (id, topic, time, eligible_for_retry, metadata) - `event.data`: Include full event with payload data - `response_data`: Include response body and headers - - name: sort_by + - name: order_by in: query required: false schema: type: string - enum: [delivery_time, event_time] - default: delivery_time - description: Sort results by delivery time or event time. - - name: sort_order + enum: [time] + default: time + description: Field to sort by. + - name: dir in: query required: false schema: type: string enum: [asc, desc] default: desc - description: Sort order (ascending or descending). + description: Sort direction. responses: "200": description: A paginated list of deliveries. content: application/json: schema: - $ref: "#/components/schemas/DeliveryListResponse" + $ref: "#/components/schemas/DeliveryPaginatedResult" examples: AdminDeliveriesListExample: value: - data: + models: - id: "del_123" status: "success" delivered_at: "2024-01-01T00:00:05Z" @@ -2572,12 +2593,16 @@ paths: attempt: 2 event: "evt_789" destination: "des_789" - next: "MTcwNDA2NzIwMA==" - prev: "" + pagination: + order_by: "time" + dir: "desc" + limit: 100 + next: "MTcwNDA2NzIwMA==" + prev: null AdminDeliveriesWithIncludeExample: summary: Response with include=event value: - data: + models: - id: "del_123" status: "success" delivered_at: "2024-01-01T00:00:05Z" @@ -2590,12 +2615,20 @@ paths: eligible_for_retry: false metadata: { "source": "crm" } destination: "des_456" - next: "" - prev: "" + pagination: + order_by: "time" + dir: "desc" + limit: 100 + next: null + prev: null "401": description: Unauthorized (Admin API Key missing or invalid). "422": description: Validation error (invalid query parameters). + content: + application/json: + schema: + $ref: "#/components/schemas/APIErrorResponse" /tenants/{tenant_id}/portal: parameters: @@ -3318,34 +3351,20 @@ paths: items: type: string description: Filter deliveries by event topic(s). Can be specified multiple times or comma-separated. - - name: start - in: query - required: false - schema: - type: string - format: date-time - description: Filter deliveries with delivered_at >= start (RFC3339 format). - - name: end - in: query - required: false - schema: - type: string - format: date-time - description: Filter deliveries with delivered_at <= end (RFC3339 format). - - name: event_start + - name: time[gte] in: query required: false schema: type: string format: date-time - description: Filter deliveries by event time >= event_start (RFC3339 format). - - name: event_end + description: Filter deliveries by event time >= value (RFC3339 or YYYY-MM-DD format). + - name: time[lte] in: query required: false schema: type: string format: date-time - description: Filter deliveries by event time <= event_end (RFC3339 format). + description: Filter deliveries by event time <= value (RFC3339 or YYYY-MM-DD format). - name: limit in: query required: false @@ -3381,33 +3400,33 @@ paths: - `event`: Include event summary (id, topic, time, eligible_for_retry, metadata) - `event.data`: Include full event with payload data - `response_data`: Include response body and headers - - name: sort_by + - name: order_by in: query required: false schema: type: string - enum: [delivery_time, event_time] - default: delivery_time - description: Sort results by delivery time or event time. - - name: sort_order + enum: [time] + default: time + description: Field to sort by. + - name: dir in: query required: false schema: type: string enum: [asc, desc] default: desc - description: Sort order (ascending or descending). + description: Sort direction. responses: "200": description: A paginated list of deliveries. content: application/json: schema: - $ref: "#/components/schemas/DeliveryListResponse" + $ref: "#/components/schemas/DeliveryPaginatedResult" examples: DeliveriesListExample: value: - data: + models: - id: "del_123" status: "success" delivered_at: "2024-01-01T00:00:05Z" @@ -3422,12 +3441,16 @@ paths: attempt: 2 event: "evt_789" destination: "des_456" - next: "MTcwNDA2NzIwMA==" - prev: "" + pagination: + order_by: "time" + dir: "desc" + limit: 100 + next: "MTcwNDA2NzIwMA==" + prev: null DeliveriesWithIncludeExample: summary: Response with include=event value: - data: + models: - id: "del_123" status: "success" delivered_at: "2024-01-01T00:00:05Z" @@ -3440,12 +3463,20 @@ paths: eligible_for_retry: false metadata: { "source": "crm" } destination: "des_456" - next: "" - prev: "" + pagination: + order_by: "time" + dir: "desc" + limit: 100 + next: null + prev: null "404": description: Tenant not found. "422": description: Validation error (invalid query parameters). + content: + application/json: + schema: + $ref: "#/components/schemas/APIErrorResponse" /tenants/{tenant_id}/deliveries/{delivery_id}: parameters: @@ -3565,7 +3596,7 @@ paths: get: tags: [Events] summary: List Events - description: Retrieves a list of events for the tenant, supporting cursor navigation (details TBD) and filtering. + description: Retrieves a list of events for the tenant, supporting cursor navigation and filtering. operationId: listTenantEvents parameters: - name: destination_id @@ -3590,13 +3621,13 @@ paths: required: false schema: type: string - description: Cursor for next page of results + description: Cursor for next page of results. - name: prev in: query required: false schema: type: string - description: Cursor for previous page of results + description: Cursor for previous page of results. - name: limit in: query required: false @@ -3605,40 +3636,48 @@ paths: default: 100 minimum: 1 maximum: 1000 - description: Number of items per page (default 100, max 1000) - - name: start + description: Number of items per page (default 100, max 1000). + - name: time[gte] in: query required: false schema: type: string format: date-time - description: Start time filter (RFC3339 format) - - name: end + description: Filter events with time >= value (RFC3339 or YYYY-MM-DD format). + - name: time[lte] in: query required: false schema: type: string format: date-time - description: End time filter (RFC3339 format) + description: Filter events with time <= value (RFC3339 or YYYY-MM-DD format). + - name: order_by + in: query + required: false + schema: + type: string + enum: [time] + default: time + description: Field to sort by. + - name: dir + in: query + required: false + schema: + type: string + enum: [asc, desc] + default: desc + description: Sort direction. responses: "200": description: A paginated list of events. content: application/json: schema: - allOf: - - $ref: "#/components/schemas/PaginatedResponse" - - type: object - properties: - data: - type: array - items: - $ref: "#/components/schemas/Event" + $ref: "#/components/schemas/EventPaginatedResult" examples: EventsListExample: value: - count: 2 - data: + models: - id: "evt_123" destination_id: "des_456" topic: "user.created" @@ -3653,11 +3692,20 @@ paths: successful_at: null metadata: { "source": "oms" } data: { "order_id": "orderid", "tracking": "1Z..." } - next: "" - prev: "" + pagination: + order_by: "time" + dir: "desc" + limit: 100 + next: null + prev: null "404": description: Tenant not found. - # Add other error responses + "422": + description: Validation error (invalid query parameters). + content: + application/json: + schema: + $ref: "#/components/schemas/APIErrorResponse" /tenants/{tenant_id}/events/{event_id}: parameters: @@ -3778,13 +3826,13 @@ paths: required: false schema: type: string - description: Cursor for next page of results + description: Cursor for next page of results. - name: prev in: query required: false schema: type: string - description: Cursor for previous page of results + description: Cursor for previous page of results. - name: limit in: query required: false @@ -3793,40 +3841,48 @@ paths: default: 100 minimum: 1 maximum: 1000 - description: Number of items per page (default 100, max 1000) - - name: start + description: Number of items per page (default 100, max 1000). + - name: time[gte] in: query required: false schema: type: string format: date-time - description: Start time filter (RFC3339 format) - - name: end + description: Filter events with time >= value (RFC3339 or YYYY-MM-DD format). + - name: time[lte] in: query required: false schema: type: string format: date-time - description: End time filter (RFC3339 format) + description: Filter events with time <= value (RFC3339 or YYYY-MM-DD format). + - name: order_by + in: query + required: false + schema: + type: string + enum: [time] + default: time + description: Field to sort by. + - name: dir + in: query + required: false + schema: + type: string + enum: [asc, desc] + default: desc + description: Sort direction. responses: "200": description: A paginated list of events for the destination. content: application/json: schema: - allOf: - - $ref: "#/components/schemas/PaginatedResponse" - - type: object - properties: - data: - type: array - items: - $ref: "#/components/schemas/Event" + $ref: "#/components/schemas/EventPaginatedResult" examples: - EventsListExample: # Same as /{tenant_id}/events example + EventsListExample: value: - count: 2 - data: + models: - id: "evt_123" destination_id: "des_456" topic: "user.created" @@ -3841,8 +3897,12 @@ paths: successful_at: null metadata: { "source": "oms" } data: { "order_id": "orderid", "tracking": "1Z..." } - next: "" - prev: "" + pagination: + order_by: "time" + dir: "desc" + limit: 100 + next: null + prev: null "404": description: Tenant or Destination not found. diff --git a/internal/apirouter/errorhandler_middleware.go b/internal/apirouter/errorhandler_middleware.go index 2bf749012..ff006ced5 100644 --- a/internal/apirouter/errorhandler_middleware.go +++ b/internal/apirouter/errorhandler_middleware.go @@ -6,6 +6,7 @@ import ( "fmt" "io" "net/http" + "strings" "github.com/gin-gonic/gin" "github.com/go-playground/validator/v10" @@ -31,6 +32,7 @@ func ErrorHandlerMiddleware() gin.HandlerFunc { type ErrorResponse struct { Err error `json:"-"` Code int `json:"-"` + Status int `json:"status"` Message string `json:"message"` Data interface{} `json:"data,omitempty"` } @@ -48,13 +50,13 @@ func (e *ErrorResponse) Parse(err error) { } if validationErrors, ok := err.(validator.ValidationErrors); ok { - out := map[string]string{} + var messages []string for _, err := range validationErrors { - out[err.Field()] = err.Tag() + messages = append(messages, formatValidationError(err.Field(), err.Tag(), err.Param())) } - e.Code = -1 + e.Code = http.StatusUnprocessableEntity e.Message = "validation error" - e.Data = out + e.Data = messages e.Err = err return } @@ -68,13 +70,13 @@ func (e *ErrorResponse) Parse(err error) { // Handle destregistry.ErrDestinationValidation var validationErr *destregistry.ErrDestinationValidation if errors.As(err, &validationErr) { - validationDetails := make(map[string]string) + var messages []string for _, detail := range validationErr.Errors { - validationDetails[detail.Field] = detail.Type + messages = append(messages, formatValidationError(detail.Field, detail.Type, "")) } e.Code = http.StatusUnprocessableEntity e.Message = "validation error" - e.Data = validationDetails + e.Data = messages e.Err = err return } @@ -83,6 +85,44 @@ func (e *ErrorResponse) Parse(err error) { e.Err = err } +// formatValidationError converts a validation error into a human-readable message. +// field is the field name, tag is the validation rule (e.g., "required", "min"), +// and param is the rule parameter (e.g., "6" for min=6). +func formatValidationError(field, tag, param string) string { + // Convert field name to lowercase for consistency + field = strings.ToLower(field) + + switch tag { + case "required": + return fmt.Sprintf("%s is required", field) + case "min": + return fmt.Sprintf("%s must be at least %s characters", field, param) + case "max": + return fmt.Sprintf("%s must be at most %s characters", field, param) + case "email": + return fmt.Sprintf("%s must be a valid email address", field) + case "url": + return fmt.Sprintf("%s must be a valid URL", field) + case "oneof": + return fmt.Sprintf("%s must be one of: %s", field, param) + case "uuid": + return fmt.Sprintf("%s must be a valid UUID", field) + case "gt": + return fmt.Sprintf("%s must be greater than %s", field, param) + case "gte": + return fmt.Sprintf("%s must be greater than or equal to %s", field, param) + case "lt": + return fmt.Sprintf("%s must be less than %s", field, param) + case "lte": + return fmt.Sprintf("%s must be less than or equal to %s", field, param) + default: + if param != "" { + return fmt.Sprintf("%s failed %s=%s validation", field, tag, param) + } + return fmt.Sprintf("%s failed %s validation", field, tag) + } +} + func isInvalidJSON(err error) bool { var syntaxError *json.SyntaxError var unmarshalTypeError *json.UnmarshalTypeError @@ -93,6 +133,7 @@ func isInvalidJSON(err error) bool { } func handleErrorResponse(c *gin.Context, response ErrorResponse) { + response.Status = response.Code c.JSON(response.Code, response) } diff --git a/internal/apirouter/errorhandler_middleware_test.go b/internal/apirouter/errorhandler_middleware_test.go new file mode 100644 index 000000000..6fb44f200 --- /dev/null +++ b/internal/apirouter/errorhandler_middleware_test.go @@ -0,0 +1,205 @@ +package apirouter_test + +import ( + "encoding/json" + "net/http" + "net/http/httptest" + "strings" + "testing" + + "github.com/gin-gonic/gin" + "github.com/go-playground/validator/v10" + "github.com/hookdeck/outpost/internal/apirouter" + "github.com/hookdeck/outpost/internal/destregistry" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func init() { + gin.SetMode(gin.TestMode) +} + +func TestErrorResponse_Parse_ValidationErrors(t *testing.T) { + t.Parallel() + + type testInput struct { + Email string `validate:"required,email"` + Password string `validate:"required,min=6"` + } + + validate := validator.New() + + t.Run("produces array of human-readable messages for validator.ValidationErrors", func(t *testing.T) { + t.Parallel() + + // Trigger validation errors + input := testInput{Email: "", Password: ""} + err := validate.Struct(input) + require.Error(t, err) + + var errorResponse apirouter.ErrorResponse + errorResponse.Parse(err) + + assert.Equal(t, http.StatusUnprocessableEntity, errorResponse.Code) + assert.Equal(t, "validation error", errorResponse.Message) + + // Data should be []string + messages, ok := errorResponse.Data.([]string) + require.True(t, ok, "Data should be []string, got %T", errorResponse.Data) + assert.Len(t, messages, 2) + + // Check that messages are human-readable (order may vary) + assert.Contains(t, messages, "email is required") + assert.Contains(t, messages, "password is required") + }) + + t.Run("includes validation param in message", func(t *testing.T) { + t.Parallel() + + // Trigger min validation error with param + input := testInput{Email: "test@example.com", Password: "abc"} + err := validate.Struct(input) + require.Error(t, err) + + var errorResponse apirouter.ErrorResponse + errorResponse.Parse(err) + + messages, ok := errorResponse.Data.([]string) + require.True(t, ok) + assert.Contains(t, messages, "password must be at least 6 characters") + }) +} + +func TestErrorResponse_Parse_DestRegistryValidation(t *testing.T) { + t.Parallel() + + t.Run("produces array of human-readable messages for destregistry.ErrDestinationValidation", func(t *testing.T) { + t.Parallel() + + err := destregistry.NewErrDestinationValidation([]destregistry.ValidationErrorDetail{ + {Field: "config.url", Type: "required"}, + {Field: "type", Type: "invalid_type"}, + }) + + var errorResponse apirouter.ErrorResponse + errorResponse.Parse(err) + + assert.Equal(t, http.StatusUnprocessableEntity, errorResponse.Code) + assert.Equal(t, "validation error", errorResponse.Message) + + messages, ok := errorResponse.Data.([]string) + require.True(t, ok, "Data should be []string, got %T", errorResponse.Data) + assert.Len(t, messages, 2) + assert.Contains(t, messages, "config.url is required") + assert.Contains(t, messages, "type failed invalid_type validation") + }) +} + +func TestHandleErrorResponse_SetsHandledAndStatus(t *testing.T) { + t.Parallel() + + router := gin.New() + router.Use(apirouter.ErrorHandlerMiddleware()) + router.GET("/test", func(c *gin.Context) { + c.Error(apirouter.NewErrBadRequest(assert.AnError)) + }) + + w := httptest.NewRecorder() + req, _ := http.NewRequest("GET", "/test", nil) + router.ServeHTTP(w, req) + + assert.Equal(t, http.StatusBadRequest, w.Code) + + var response map[string]interface{} + err := json.Unmarshal(w.Body.Bytes(), &response) + require.NoError(t, err) + + assert.Equal(t, float64(http.StatusBadRequest), response["status"]) + assert.Equal(t, assert.AnError.Error(), response["message"]) +} + +func TestHandleErrorResponse_ValidationErrorFormat(t *testing.T) { + t.Parallel() + + type requestBody struct { + Name string `json:"name" binding:"required"` + } + + router := gin.New() + router.Use(apirouter.ErrorHandlerMiddleware()) + router.POST("/test", func(c *gin.Context) { + var body requestBody + if err := c.ShouldBindJSON(&body); err != nil { + apirouter.AbortWithValidationError(c, err) + return + } + c.JSON(http.StatusOK, body) + }) + + w := httptest.NewRecorder() + // Use empty JSON object to trigger validation error (not JSON parse error) + req, _ := http.NewRequest("POST", "/test", strings.NewReader("{}")) + req.Header.Set("Content-Type", "application/json") + router.ServeHTTP(w, req) + + assert.Equal(t, http.StatusUnprocessableEntity, w.Code) + + var response map[string]interface{} + err := json.Unmarshal(w.Body.Bytes(), &response) + require.NoError(t, err) + + assert.Equal(t, float64(http.StatusUnprocessableEntity), response["status"]) + assert.Equal(t, "validation error", response["message"]) + + // Data should be an array + data, ok := response["data"].([]interface{}) + require.True(t, ok, "data should be an array, got %T", response["data"]) + assert.Len(t, data, 1) + assert.Equal(t, "name is required", data[0]) +} + +func TestErrorResponse_NotFoundFormat(t *testing.T) { + t.Parallel() + + router := gin.New() + router.Use(apirouter.ErrorHandlerMiddleware()) + router.GET("/test", func(c *gin.Context) { + c.Error(apirouter.NewErrNotFound("tenant")) + }) + + w := httptest.NewRecorder() + req, _ := http.NewRequest("GET", "/test", nil) + router.ServeHTTP(w, req) + + assert.Equal(t, http.StatusNotFound, w.Code) + + var response map[string]interface{} + err := json.Unmarshal(w.Body.Bytes(), &response) + require.NoError(t, err) + + assert.Equal(t, float64(http.StatusNotFound), response["status"]) + assert.Equal(t, "tenant not found", response["message"]) +} + +func TestErrorResponse_InternalServerErrorFormat(t *testing.T) { + t.Parallel() + + router := gin.New() + router.Use(apirouter.ErrorHandlerMiddleware()) + router.GET("/test", func(c *gin.Context) { + c.Error(apirouter.NewErrInternalServer(assert.AnError)) + }) + + w := httptest.NewRecorder() + req, _ := http.NewRequest("GET", "/test", nil) + router.ServeHTTP(w, req) + + assert.Equal(t, http.StatusInternalServerError, w.Code) + + var response map[string]interface{} + err := json.Unmarshal(w.Body.Bytes(), &response) + require.NoError(t, err) + + assert.Equal(t, float64(http.StatusInternalServerError), response["status"]) + assert.Equal(t, "internal server error", response["message"]) +} diff --git a/internal/apirouter/legacy_handlers.go b/internal/apirouter/legacy_handlers.go index 1e45ede7d..18195f5ee 100644 --- a/internal/apirouter/legacy_handlers.go +++ b/internal/apirouter/legacy_handlers.go @@ -121,6 +121,13 @@ func (h *LegacyHandlers) ListEventsByDestination(c *gin.Context) { } destinationID := c.Param("destinationID") + // Parse and validate cursors (next/prev are mutually exclusive) + cursors, errResp := ParseCursors(c) + if errResp != nil { + AbortWithError(c, errResp.Code, *errResp) + return + } + // Parse pagination params limit := 100 if limitStr := c.Query("limit"); limitStr != "" { @@ -134,8 +141,8 @@ func (h *LegacyHandlers) ListEventsByDestination(c *gin.Context) { TenantID: tenant.ID, DestinationIDs: []string{destinationID}, Limit: limit, - Next: c.Query("next"), - Prev: c.Query("prev"), + Next: cursors.Next, + Prev: cursors.Prev, SortOrder: "desc", }) if err != nil { diff --git a/internal/apirouter/log_handlers.go b/internal/apirouter/log_handlers.go index d0ee22d5e..7e2ee0b21 100644 --- a/internal/apirouter/log_handlers.go +++ b/internal/apirouter/log_handlers.go @@ -131,18 +131,16 @@ type APIEvent struct { Data map[string]interface{} `json:"data,omitempty"` } -// ListDeliveriesResponse is the response for ListDeliveries -type ListDeliveriesResponse struct { - Data []APIDelivery `json:"data"` - Next string `json:"next,omitempty"` - Prev string `json:"prev,omitempty"` +// DeliveryPaginatedResult is the paginated response for listing deliveries. +type DeliveryPaginatedResult struct { + Models []APIDelivery `json:"models"` + Pagination SeekPagination `json:"pagination"` } -// ListEventsResponse is the response for ListEvents -type ListEventsResponse struct { - Data []APIEvent `json:"data"` - Next string `json:"next,omitempty"` - Prev string `json:"prev,omitempty"` +// EventPaginatedResult is the paginated response for listing events. +type EventPaginatedResult struct { + Models []APIEvent `json:"models"` + Pagination SeekPagination `json:"pagination"` } // toAPIDelivery converts a DeliveryEvent to APIDelivery with expand options @@ -203,34 +201,40 @@ func (h *LogHandlers) ListDeliveries(c *gin.Context) { } func (h *LogHandlers) listDeliveriesInternal(c *gin.Context, tenantID string) { - var start, end *time.Time - if startStr := c.Query("start"); startStr != "" { - t, err := time.Parse(time.RFC3339, startStr) - if err != nil { - AbortWithError(c, http.StatusUnprocessableEntity, ErrorResponse{ - Code: http.StatusUnprocessableEntity, - Message: "validation error", - Data: map[string]string{ - "query.start": "invalid format, expected RFC3339", - }, - }) - return - } - start = &t - } - if endStr := c.Query("end"); endStr != "" { - t, err := time.Parse(time.RFC3339, endStr) - if err != nil { - AbortWithError(c, http.StatusUnprocessableEntity, ErrorResponse{ - Code: http.StatusUnprocessableEntity, - Message: "validation error", - Data: map[string]string{ - "query.end": "invalid format, expected RFC3339", - }, - }) - return - } - end = &t + // Parse and validate cursors (next/prev are mutually exclusive) + cursors, errResp := ParseCursors(c) + if errResp != nil { + AbortWithError(c, errResp.Code, *errResp) + return + } + + // Parse and validate dir (sort direction) + dir, errResp := ParseDir(c) + if errResp != nil { + AbortWithError(c, errResp.Code, *errResp) + return + } + if dir == "" { + dir = "desc" + } + + // Parse and validate order_by (time only) + orderBy, errResp := ParseOrderBy(c, []string{"time"}) + if errResp != nil { + AbortWithError(c, errResp.Code, *errResp) + return + } + if orderBy == "" { + orderBy = "time" + } + // Note: order_by is informational only for now - store always sorts by time + _ = orderBy + + // Parse time date filters + deliveryTimeFilter, errResp := ParseDateFilter(c, "time") + if errResp != nil { + AbortWithError(c, errResp.Code, *errResp) + return } limit := parseLimit(c, 100, 1000) @@ -240,33 +244,22 @@ func (h *LogHandlers) listDeliveriesInternal(c *gin.Context, tenantID string) { destinationIDs = []string{destID} } - sortOrder := c.Query("sort_order") - if sortOrder == "" { - sortOrder = "desc" - } - if sortOrder != "asc" && sortOrder != "desc" { - AbortWithError(c, http.StatusUnprocessableEntity, ErrorResponse{ - Code: http.StatusUnprocessableEntity, - Message: "validation error", - Data: map[string]string{ - "query.sort_order": "must be 'asc' or 'desc'", - }, - }) - return - } - req := logstore.ListDeliveryEventRequest{ TenantID: tenantID, EventID: c.Query("event_id"), DestinationIDs: destinationIDs, Status: c.Query("status"), Topics: parseQueryArray(c, "topic"), - Start: start, - End: end, - Limit: limit, - Next: c.Query("next"), - Prev: c.Query("prev"), - SortOrder: sortOrder, + TimeFilter: logstore.TimeFilter{ + GTE: deliveryTimeFilter.GTE, + LTE: deliveryTimeFilter.LTE, + GT: deliveryTimeFilter.GT, + LT: deliveryTimeFilter.LT, + }, + Limit: limit, + Next: cursors.Next, + Prev: cursors.Prev, + SortOrder: dir, } response, err := h.logStore.ListDeliveryEvent(c.Request.Context(), req) @@ -286,10 +279,15 @@ func (h *LogHandlers) listDeliveriesInternal(c *gin.Context, tenantID string) { apiDeliveries[i] = toAPIDelivery(de, includeOpts) } - c.JSON(http.StatusOK, ListDeliveriesResponse{ - Data: apiDeliveries, - Next: response.Next, - Prev: response.Prev, + c.JSON(http.StatusOK, DeliveryPaginatedResult{ + Models: apiDeliveries, + Pagination: SeekPagination{ + OrderBy: orderBy, + Dir: dir, + Limit: limit, + Next: CursorToPtr(response.Next), + Prev: CursorToPtr(response.Prev), + }, }) } @@ -371,34 +369,40 @@ func (h *LogHandlers) ListEvents(c *gin.Context) { } func (h *LogHandlers) listEventsInternal(c *gin.Context, tenantID string) { - var start, end *time.Time - if startStr := c.Query("start"); startStr != "" { - t, err := time.Parse(time.RFC3339, startStr) - if err != nil { - AbortWithError(c, http.StatusUnprocessableEntity, ErrorResponse{ - Code: http.StatusUnprocessableEntity, - Message: "validation error", - Data: map[string]string{ - "query.start": "invalid format, expected RFC3339", - }, - }) - return - } - start = &t - } - if endStr := c.Query("end"); endStr != "" { - t, err := time.Parse(time.RFC3339, endStr) - if err != nil { - AbortWithError(c, http.StatusUnprocessableEntity, ErrorResponse{ - Code: http.StatusUnprocessableEntity, - Message: "validation error", - Data: map[string]string{ - "query.end": "invalid format, expected RFC3339", - }, - }) - return - } - end = &t + // Parse and validate cursors (next/prev are mutually exclusive) + cursors, errResp := ParseCursors(c) + if errResp != nil { + AbortWithError(c, errResp.Code, *errResp) + return + } + + // Parse and validate dir (sort direction) + dir, errResp := ParseDir(c) + if errResp != nil { + AbortWithError(c, errResp.Code, *errResp) + return + } + if dir == "" { + dir = "desc" + } + + // Parse and validate order_by (time only) + orderBy, errResp := ParseOrderBy(c, []string{"time"}) + if errResp != nil { + AbortWithError(c, errResp.Code, *errResp) + return + } + if orderBy == "" { + orderBy = "time" + } + // Note: order_by is informational only for now - store always sorts by time + _ = orderBy + + // Parse time date filters + eventTimeFilter, errResp := ParseDateFilter(c, "time") + if errResp != nil { + AbortWithError(c, errResp.Code, *errResp) + return } limit := parseLimit(c, 100, 1000) @@ -408,31 +412,20 @@ func (h *LogHandlers) listEventsInternal(c *gin.Context, tenantID string) { destinationIDs = []string{destID} } - sortOrder := c.Query("sort_order") - if sortOrder == "" { - sortOrder = "desc" - } - if sortOrder != "asc" && sortOrder != "desc" { - AbortWithError(c, http.StatusUnprocessableEntity, ErrorResponse{ - Code: http.StatusUnprocessableEntity, - Message: "validation error", - Data: map[string]string{ - "query.sort_order": "must be 'asc' or 'desc'", - }, - }) - return - } - req := logstore.ListEventRequest{ TenantID: tenantID, DestinationIDs: destinationIDs, Topics: parseQueryArray(c, "topic"), - EventStart: start, - EventEnd: end, - Limit: limit, - Next: c.Query("next"), - Prev: c.Query("prev"), - SortOrder: sortOrder, + TimeFilter: logstore.TimeFilter{ + GTE: eventTimeFilter.GTE, + LTE: eventTimeFilter.LTE, + GT: eventTimeFilter.GT, + LT: eventTimeFilter.LT, + }, + Limit: limit, + Next: cursors.Next, + Prev: cursors.Prev, + SortOrder: dir, } response, err := h.logStore.ListEvent(c.Request.Context(), req) @@ -457,9 +450,14 @@ func (h *LogHandlers) listEventsInternal(c *gin.Context, tenantID string) { } } - c.JSON(http.StatusOK, ListEventsResponse{ - Data: apiEvents, - Next: response.Next, - Prev: response.Prev, + c.JSON(http.StatusOK, EventPaginatedResult{ + Models: apiEvents, + Pagination: SeekPagination{ + OrderBy: orderBy, + Dir: dir, + Limit: limit, + Next: CursorToPtr(response.Next), + Prev: CursorToPtr(response.Prev), + }, }) } diff --git a/internal/apirouter/log_handlers_test.go b/internal/apirouter/log_handlers_test.go index f80e2c9a5..3bfc69acb 100644 --- a/internal/apirouter/log_handlers_test.go +++ b/internal/apirouter/log_handlers_test.go @@ -46,7 +46,7 @@ func TestListDeliveries(t *testing.T) { var response map[string]interface{} require.NoError(t, json.Unmarshal(w.Body.Bytes(), &response)) - data := response["data"].([]interface{}) + data := response["models"].([]interface{}) assert.Len(t, data, 0) }) @@ -91,7 +91,7 @@ func TestListDeliveries(t *testing.T) { var response map[string]interface{} require.NoError(t, json.Unmarshal(w.Body.Bytes(), &response)) - data := response["data"].([]interface{}) + data := response["models"].([]interface{}) assert.Len(t, data, 1) firstDelivery := data[0].(map[string]interface{}) @@ -111,7 +111,7 @@ func TestListDeliveries(t *testing.T) { var response map[string]interface{} require.NoError(t, json.Unmarshal(w.Body.Bytes(), &response)) - data := response["data"].([]interface{}) + data := response["models"].([]interface{}) require.Len(t, data, 1) firstDelivery := data[0].(map[string]interface{}) @@ -132,7 +132,7 @@ func TestListDeliveries(t *testing.T) { var response map[string]interface{} require.NoError(t, json.Unmarshal(w.Body.Bytes(), &response)) - data := response["data"].([]interface{}) + data := response["models"].([]interface{}) require.Len(t, data, 1) firstDelivery := data[0].(map[string]interface{}) @@ -151,7 +151,7 @@ func TestListDeliveries(t *testing.T) { var response map[string]interface{} require.NoError(t, json.Unmarshal(w.Body.Bytes(), &response)) - data := response["data"].([]interface{}) + data := response["models"].([]interface{}) assert.Len(t, data, 1) }) @@ -165,7 +165,7 @@ func TestListDeliveries(t *testing.T) { var response map[string]interface{} require.NoError(t, json.Unmarshal(w.Body.Bytes(), &response)) - data := response["data"].([]interface{}) + data := response["models"].([]interface{}) assert.Len(t, data, 0) }) @@ -187,7 +187,7 @@ func TestListDeliveries(t *testing.T) { var response map[string]interface{} require.NoError(t, json.Unmarshal(w.Body.Bytes(), &response)) - data := response["data"].([]interface{}) + data := response["models"].([]interface{}) require.Len(t, data, 1) firstDelivery := data[0].(map[string]interface{}) @@ -239,7 +239,7 @@ func TestListDeliveries(t *testing.T) { var response map[string]interface{} require.NoError(t, json.Unmarshal(w.Body.Bytes(), &response)) - data := response["data"].([]interface{}) + data := response["models"].([]interface{}) // Find the delivery we just created var foundDelivery map[string]interface{} for _, d := range data { @@ -266,7 +266,7 @@ func TestListDeliveries(t *testing.T) { var response map[string]interface{} require.NoError(t, json.Unmarshal(w.Body.Bytes(), &response)) - data := response["data"].([]interface{}) + data := response["models"].([]interface{}) require.GreaterOrEqual(t, len(data), 1) firstDelivery := data[0].(map[string]interface{}) @@ -276,17 +276,17 @@ func TestListDeliveries(t *testing.T) { assert.NotNil(t, event["topic"]) }) - t.Run("should return validation error for invalid sort_order", func(t *testing.T) { + t.Run("should return validation error for invalid dir", func(t *testing.T) { w := httptest.NewRecorder() - req, _ := http.NewRequest("GET", baseAPIPath+"/tenants/"+tenantID+"/deliveries?sort_order=invalid", nil) + req, _ := http.NewRequest("GET", baseAPIPath+"/tenants/"+tenantID+"/deliveries?dir=invalid", nil) result.router.ServeHTTP(w, req) assert.Equal(t, http.StatusUnprocessableEntity, w.Code) }) - t.Run("should accept valid sort_order param", func(t *testing.T) { + t.Run("should accept valid dir param", func(t *testing.T) { w := httptest.NewRecorder() - req, _ := http.NewRequest("GET", baseAPIPath+"/tenants/"+tenantID+"/deliveries?sort_order=asc", nil) + req, _ := http.NewRequest("GET", baseAPIPath+"/tenants/"+tenantID+"/deliveries?dir=asc", nil) result.router.ServeHTTP(w, req) assert.Equal(t, http.StatusOK, w.Code) @@ -540,7 +540,7 @@ func TestListEvents(t *testing.T) { var response map[string]interface{} require.NoError(t, json.Unmarshal(w.Body.Bytes(), &response)) - data := response["data"].([]interface{}) + data := response["models"].([]interface{}) assert.Len(t, data, 0) }) @@ -588,7 +588,7 @@ func TestListEvents(t *testing.T) { var response map[string]interface{} require.NoError(t, json.Unmarshal(w.Body.Bytes(), &response)) - data := response["data"].([]interface{}) + data := response["models"].([]interface{}) assert.Len(t, data, 1) firstEvent := data[0].(map[string]interface{}) @@ -607,7 +607,7 @@ func TestListEvents(t *testing.T) { var response map[string]interface{} require.NoError(t, json.Unmarshal(w.Body.Bytes(), &response)) - data := response["data"].([]interface{}) + data := response["models"].([]interface{}) assert.GreaterOrEqual(t, len(data), 1) }) @@ -621,7 +621,7 @@ func TestListEvents(t *testing.T) { var response map[string]interface{} require.NoError(t, json.Unmarshal(w.Body.Bytes(), &response)) - data := response["data"].([]interface{}) + data := response["models"].([]interface{}) assert.Len(t, data, 0) }) @@ -635,7 +635,7 @@ func TestListEvents(t *testing.T) { var response map[string]interface{} require.NoError(t, json.Unmarshal(w.Body.Bytes(), &response)) - data := response["data"].([]interface{}) + data := response["models"].([]interface{}) assert.GreaterOrEqual(t, len(data), 1) for _, item := range data { event := item.(map[string]interface{}) @@ -651,33 +651,33 @@ func TestListEvents(t *testing.T) { assert.Equal(t, http.StatusNotFound, w.Code) }) - t.Run("should return validation error for invalid start time", func(t *testing.T) { + t.Run("should return validation error for invalid time filter", func(t *testing.T) { w := httptest.NewRecorder() - req, _ := http.NewRequest("GET", baseAPIPath+"/tenants/"+tenantID+"/events?start=invalid", nil) + req, _ := http.NewRequest("GET", baseAPIPath+"/tenants/"+tenantID+"/events?time[gte]=invalid", nil) result.router.ServeHTTP(w, req) assert.Equal(t, http.StatusUnprocessableEntity, w.Code) }) - t.Run("should return validation error for invalid end time", func(t *testing.T) { + t.Run("should return validation error for invalid time lte filter", func(t *testing.T) { w := httptest.NewRecorder() - req, _ := http.NewRequest("GET", baseAPIPath+"/tenants/"+tenantID+"/events?end=invalid", nil) + req, _ := http.NewRequest("GET", baseAPIPath+"/tenants/"+tenantID+"/events?time[lte]=invalid", nil) result.router.ServeHTTP(w, req) assert.Equal(t, http.StatusUnprocessableEntity, w.Code) }) - t.Run("should return validation error for invalid sort_order", func(t *testing.T) { + t.Run("should return validation error for invalid dir", func(t *testing.T) { w := httptest.NewRecorder() - req, _ := http.NewRequest("GET", baseAPIPath+"/tenants/"+tenantID+"/events?sort_order=invalid", nil) + req, _ := http.NewRequest("GET", baseAPIPath+"/tenants/"+tenantID+"/events?dir=invalid", nil) result.router.ServeHTTP(w, req) assert.Equal(t, http.StatusUnprocessableEntity, w.Code) }) - t.Run("should accept valid sort_order param", func(t *testing.T) { + t.Run("should accept valid dir param", func(t *testing.T) { w := httptest.NewRecorder() - req, _ := http.NewRequest("GET", baseAPIPath+"/tenants/"+tenantID+"/events?sort_order=asc", nil) + req, _ := http.NewRequest("GET", baseAPIPath+"/tenants/"+tenantID+"/events?dir=asc", nil) result.router.ServeHTTP(w, req) assert.Equal(t, http.StatusOK, w.Code) diff --git a/internal/apirouter/pagination.go b/internal/apirouter/pagination.go new file mode 100644 index 000000000..aa96a7b04 --- /dev/null +++ b/internal/apirouter/pagination.go @@ -0,0 +1,258 @@ +package apirouter + +import ( + "fmt" + "net/http" + "slices" + "sort" + "strconv" + "strings" + "time" + + "github.com/gin-gonic/gin" +) + +// SeekPagination represents cursor-based pagination metadata for list responses. +type SeekPagination struct { + OrderBy string `json:"order_by"` + Dir string `json:"dir"` + Limit int `json:"limit"` + Next *string `json:"next"` + Prev *string `json:"prev"` +} + +// CursorToPtr converts an empty cursor string to nil, or returns a pointer to the string. +// Returns null for empty cursors instead of empty string. +func CursorToPtr(cursor string) *string { + if cursor == "" { + return nil + } + return &cursor +} + +// CursorParams holds the parsed cursor values for pagination. +type CursorParams struct { + Next string + Prev string +} + +// ParseCursors parses the "next" and "prev" query parameters. +// Returns an error if both are provided (mutually exclusive). +func ParseCursors(c *gin.Context) (CursorParams, *ErrorResponse) { + next := c.Query("next") + prev := c.Query("prev") + + if next != "" && prev != "" { + return CursorParams{}, &ErrorResponse{ + Code: http.StatusBadRequest, + Message: "cannot specify both 'next' and 'prev' cursors", + } + } + + return CursorParams{Next: next, Prev: prev}, nil +} + +// ParseDir parses the "dir" query parameter for sort direction. +// Returns the direction (empty string if not provided) and any validation error response. +// Valid values are "asc" and "desc". Caller/store should apply default if empty. +func ParseDir(c *gin.Context) (string, *ErrorResponse) { + dir := c.Query("dir") + if dir == "" { + return "", nil + } + if dir != "asc" && dir != "desc" { + return "", &ErrorResponse{ + Code: http.StatusUnprocessableEntity, + Message: "validation error", + Data: map[string]string{ + "query.dir": "must be 'asc' or 'desc'", + }, + } + } + return dir, nil +} + +// ParseOrderBy parses the "order_by" query parameter and validates against allowed values. +// Returns the order_by value (empty string if not provided) and any validation error response. +// Caller/store should apply default if empty. +func ParseOrderBy(c *gin.Context, allowedValues []string) (string, *ErrorResponse) { + orderBy := c.Query("order_by") + if orderBy == "" { + return "", nil + } + + if slices.Contains(allowedValues, orderBy) { + return orderBy, nil + } + + return "", &ErrorResponse{ + Code: http.StatusUnprocessableEntity, + Message: "validation error", + Data: map[string]string{ + "query.order_by": fmt.Sprintf("must be one of: %v", allowedValues), + }, + } +} + +// ParseArrayParam parses bracket notation array parameters (e.g., field[0]=a&field[1]=b). +// Returns the parsed array in index order. Supports both numeric indices [0], [1] and +// simple repeated params field[]=a&field[]=b. +func ParseArrayParam(c *gin.Context, fieldName string) []string { + queryMap := c.Request.URL.Query() + prefix := fieldName + "[" + + // Collect all matching params + type indexedValue struct { + index int + value string + } + var values []indexedValue + var unindexed []string + + for key, vals := range queryMap { + if !strings.HasPrefix(key, prefix) { + continue + } + if len(vals) == 0 { + continue + } + + // Extract the part between [ and ] + suffix := key[len(prefix):] + if !strings.HasSuffix(suffix, "]") { + continue + } + indexStr := suffix[:len(suffix)-1] + + if indexStr == "" { + // field[]=value format - collect all values + unindexed = append(unindexed, vals...) + } else { + // field[0]=value format - parse index + idx, err := strconv.Atoi(indexStr) + if err != nil { + continue // skip non-numeric indices + } + values = append(values, indexedValue{index: idx, value: vals[0]}) + } + } + + // If we have indexed values, sort by index and return + if len(values) > 0 { + sort.Slice(values, func(i, j int) bool { + return values[i].index < values[j].index + }) + result := make([]string, len(values)) + for i, v := range values { + result[i] = v.value + } + return result + } + + // Return unindexed values if any + return unindexed +} + +// DateFilterResult holds the parsed date filter values. +type DateFilterResult struct { + GTE *time.Time // Greater than or equal (>=) + LTE *time.Time // Less than or equal (<=) + GT *time.Time // Greater than (>) + LT *time.Time // Less than (<) +} + +// unsupportedDateOps lists operators we recognize but don't support yet. +// When stores add support, move these to supported and update ParseDateFilter. +var unsupportedDateOps = []string{"any"} + +// ParseDateFilter parses bracket notation date filters (e.g., field[gte], field[lte], field[gt], field[lt]). +// Returns 400 for unsupported operators (any). +// Accepts both RFC3339 format (2024-01-01T00:00:00Z) and date-only format (2024-01-01). +func ParseDateFilter(c *gin.Context, fieldName string) (*DateFilterResult, *ErrorResponse) { + // Check for unsupported operators first + for _, op := range unsupportedDateOps { + key := fieldName + "[" + op + "]" + if c.Query(key) != "" { + return nil, &ErrorResponse{ + Code: http.StatusBadRequest, + Message: fmt.Sprintf("operator '%s' is not supported, use 'gte', 'lte', 'gt', or 'lt'", op), + } + } + } + + result := &DateFilterResult{} + + // Parse [gte] - greater than or equal + gteKey := fieldName + "[gte]" + if gteStr := c.Query(gteKey); gteStr != "" { + t, err := parseDateTime(gteStr) + if err != nil { + return nil, dateFormatError(gteKey) + } + result.GTE = &t + } + + // Parse [lte] - less than or equal + lteKey := fieldName + "[lte]" + if lteStr := c.Query(lteKey); lteStr != "" { + t, err := parseDateTime(lteStr) + if err != nil { + return nil, dateFormatError(lteKey) + } + result.LTE = &t + } + + // Parse [gt] - greater than + gtKey := fieldName + "[gt]" + if gtStr := c.Query(gtKey); gtStr != "" { + t, err := parseDateTime(gtStr) + if err != nil { + return nil, dateFormatError(gtKey) + } + result.GT = &t + } + + // Parse [lt] - less than + ltKey := fieldName + "[lt]" + if ltStr := c.Query(ltKey); ltStr != "" { + t, err := parseDateTime(ltStr) + if err != nil { + return nil, dateFormatError(ltKey) + } + result.LT = &t + } + + return result, nil +} + +// dateFormatError returns a validation error for invalid date format. +func dateFormatError(key string) *ErrorResponse { + return &ErrorResponse{ + Code: http.StatusUnprocessableEntity, + Message: "validation error", + Data: map[string]string{ + "query." + key: "invalid format, expected RFC3339 (e.g. 2024-01-15T09:30:00Z or 2024-01-15T09:30:00.123Z) or YYYY-MM-DD", + }, + } +} + +// Note: String and number filters are not yet supported by the stores. +// These types are defined for future use when stores add support. +// Attempting to use them will return 400 errors from the handlers. + +// parseDateTime attempts to parse a date string in RFC3339, RFC3339Nano, or date-only (YYYY-MM-DD) format. +func parseDateTime(s string) (time.Time, error) { + // Try RFC3339Nano first (handles milliseconds/microseconds like 2024-01-01T00:00:00.123Z) + if t, err := time.Parse(time.RFC3339Nano, s); err == nil { + return t, nil + } + // Try RFC3339 (no fractional seconds) + if t, err := time.Parse(time.RFC3339, s); err == nil { + return t, nil + } + // Try date-only format + if t, err := time.Parse("2006-01-02", s); err == nil { + return t, nil + } + return time.Time{}, fmt.Errorf("invalid date format") +} diff --git a/internal/apirouter/pagination_test.go b/internal/apirouter/pagination_test.go new file mode 100644 index 000000000..3a3b32fca --- /dev/null +++ b/internal/apirouter/pagination_test.go @@ -0,0 +1,456 @@ +package apirouter_test + +import ( + "net/http" + "net/http/httptest" + "net/url" + "testing" + "time" + + "github.com/gin-gonic/gin" + "github.com/hookdeck/outpost/internal/apirouter" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func init() { + gin.SetMode(gin.TestMode) +} + +func TestCursorToPtr(t *testing.T) { + t.Run("empty string returns nil", func(t *testing.T) { + result := apirouter.CursorToPtr("") + assert.Nil(t, result) + }) + + t.Run("non-empty string returns pointer", func(t *testing.T) { + result := apirouter.CursorToPtr("abc123") + require.NotNil(t, result) + assert.Equal(t, "abc123", *result) + }) +} + +func TestParseCursors(t *testing.T) { + tests := []struct { + name string + queryParams map[string]string + wantNext string + wantPrev string + wantErrCode int + }{ + { + name: "no cursors", + queryParams: map[string]string{}, + wantNext: "", + wantPrev: "", + }, + { + name: "next only", + queryParams: map[string]string{"next": "cursor123"}, + wantNext: "cursor123", + wantPrev: "", + }, + { + name: "prev only", + queryParams: map[string]string{"prev": "cursor456"}, + wantNext: "", + wantPrev: "cursor456", + }, + { + name: "both next and prev returns error", + queryParams: map[string]string{"next": "cursor123", "prev": "cursor456"}, + wantErrCode: http.StatusBadRequest, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + c, _ := createTestContext(tt.queryParams) + + result, errResp := apirouter.ParseCursors(c) + + if tt.wantErrCode != 0 { + require.NotNil(t, errResp) + assert.Equal(t, tt.wantErrCode, errResp.Code) + assert.Contains(t, errResp.Message, "both") + return + } + + assert.Nil(t, errResp) + assert.Equal(t, tt.wantNext, result.Next) + assert.Equal(t, tt.wantPrev, result.Prev) + }) + } +} + +func TestParseDir(t *testing.T) { + tests := []struct { + name string + queryDir string + wantDir string + wantErrCode int + }{ + { + name: "empty returns empty (no default)", + queryDir: "", + wantDir: "", + }, + { + name: "asc is valid", + queryDir: "asc", + wantDir: "asc", + }, + { + name: "desc is valid", + queryDir: "desc", + wantDir: "desc", + }, + { + name: "invalid value returns error", + queryDir: "invalid", + wantErrCode: http.StatusUnprocessableEntity, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + c, _ := createTestContext(map[string]string{"dir": tt.queryDir}) + + dir, errResp := apirouter.ParseDir(c) + + if tt.wantErrCode != 0 { + require.NotNil(t, errResp) + assert.Equal(t, tt.wantErrCode, errResp.Code) + } else { + assert.Nil(t, errResp) + assert.Equal(t, tt.wantDir, dir) + } + }) + } +} + +func TestParseOrderBy(t *testing.T) { + allowedValues := []string{"created_at", "time"} + + tests := []struct { + name string + queryOrderBy string + wantOrderBy string + wantErrCode int + }{ + { + name: "empty returns empty (no default)", + queryOrderBy: "", + wantOrderBy: "", + }, + { + name: "valid value returns value", + queryOrderBy: "time", + wantOrderBy: "time", + }, + { + name: "another valid value", + queryOrderBy: "created_at", + wantOrderBy: "created_at", + }, + { + name: "invalid value returns error", + queryOrderBy: "invalid_field", + wantErrCode: http.StatusUnprocessableEntity, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + c, _ := createTestContext(map[string]string{"order_by": tt.queryOrderBy}) + + orderBy, errResp := apirouter.ParseOrderBy(c, allowedValues) + + if tt.wantErrCode != 0 { + require.NotNil(t, errResp) + assert.Equal(t, tt.wantErrCode, errResp.Code) + } else { + assert.Nil(t, errResp) + assert.Equal(t, tt.wantOrderBy, orderBy) + } + }) + } +} + +func TestParseDateFilter(t *testing.T) { + tests := []struct { + name string + fieldName string + queryParams map[string]string + wantGTE *time.Time + wantLTE *time.Time + wantErrCode int + }{ + { + name: "no params returns empty result", + fieldName: "time", + queryParams: map[string]string{}, + }, + { + name: "RFC3339 gte param", + fieldName: "time", + queryParams: map[string]string{ + "time[gte]": "2024-01-15T10:30:00Z", + }, + wantGTE: timePtr(time.Date(2024, 1, 15, 10, 30, 0, 0, time.UTC)), + }, + { + name: "RFC3339 lte param", + fieldName: "time", + queryParams: map[string]string{ + "time[lte]": "2024-01-31T23:59:59Z", + }, + wantLTE: timePtr(time.Date(2024, 1, 31, 23, 59, 59, 0, time.UTC)), + }, + { + name: "both gte and lte params", + fieldName: "created_at", + queryParams: map[string]string{ + "created_at[gte]": "2024-01-01T00:00:00Z", + "created_at[lte]": "2024-01-31T23:59:59Z", + }, + wantGTE: timePtr(time.Date(2024, 1, 1, 0, 0, 0, 0, time.UTC)), + wantLTE: timePtr(time.Date(2024, 1, 31, 23, 59, 59, 0, time.UTC)), + }, + { + name: "date-only format", + fieldName: "time", + queryParams: map[string]string{ + "time[gte]": "2024-01-15", + }, + wantGTE: timePtr(time.Date(2024, 1, 15, 0, 0, 0, 0, time.UTC)), + }, + { + name: "invalid gte format returns error", + fieldName: "time", + queryParams: map[string]string{ + "time[gte]": "not-a-date", + }, + wantErrCode: http.StatusUnprocessableEntity, + }, + { + name: "invalid lte format returns error", + fieldName: "time", + queryParams: map[string]string{ + "time[lte]": "invalid", + }, + wantErrCode: http.StatusUnprocessableEntity, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + c, _ := createTestContext(tt.queryParams) + + result, errResp := apirouter.ParseDateFilter(c, tt.fieldName) + + if tt.wantErrCode != 0 { + require.NotNil(t, errResp) + assert.Equal(t, tt.wantErrCode, errResp.Code) + return + } + + require.Nil(t, errResp) + require.NotNil(t, result) + + assertTimePtr(t, "GTE", tt.wantGTE, result.GTE) + assertTimePtr(t, "LTE", tt.wantLTE, result.LTE) + }) + } +} + +func TestParseDateFilter_GTLTOperators(t *testing.T) { + tests := []struct { + name string + fieldName string + queryParams map[string]string + wantGT *time.Time + wantLT *time.Time + wantErrCode int + }{ + { + name: "gt param", + fieldName: "time", + queryParams: map[string]string{ + "time[gt]": "2024-01-15T10:30:00Z", + }, + wantGT: timePtr(time.Date(2024, 1, 15, 10, 30, 0, 0, time.UTC)), + }, + { + name: "lt param", + fieldName: "time", + queryParams: map[string]string{ + "time[lt]": "2024-01-31T23:59:59Z", + }, + wantLT: timePtr(time.Date(2024, 1, 31, 23, 59, 59, 0, time.UTC)), + }, + { + name: "all four operators", + fieldName: "time", + queryParams: map[string]string{ + "time[gte]": "2024-01-01T00:00:00Z", + "time[lte]": "2024-01-31T23:59:59Z", + "time[gt]": "2024-01-02T00:00:00Z", + "time[lt]": "2024-01-30T23:59:59Z", + }, + wantGT: timePtr(time.Date(2024, 1, 2, 0, 0, 0, 0, time.UTC)), + wantLT: timePtr(time.Date(2024, 1, 30, 23, 59, 59, 0, time.UTC)), + }, + { + name: "invalid gt format returns error", + fieldName: "time", + queryParams: map[string]string{ + "time[gt]": "not-a-date", + }, + wantErrCode: http.StatusUnprocessableEntity, + }, + { + name: "invalid lt format returns error", + fieldName: "time", + queryParams: map[string]string{ + "time[lt]": "invalid", + }, + wantErrCode: http.StatusUnprocessableEntity, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + c, _ := createTestContext(tt.queryParams) + + result, errResp := apirouter.ParseDateFilter(c, tt.fieldName) + + if tt.wantErrCode != 0 { + require.NotNil(t, errResp) + assert.Equal(t, tt.wantErrCode, errResp.Code) + return + } + + require.Nil(t, errResp) + require.NotNil(t, result) + + assertTimePtr(t, "GT", tt.wantGT, result.GT) + assertTimePtr(t, "LT", tt.wantLT, result.LT) + }) + } +} + +func TestParseDateFilter_UnsupportedOperators(t *testing.T) { + unsupportedOps := []string{"any"} + + for _, op := range unsupportedOps { + t.Run(op+" returns 400", func(t *testing.T) { + c, _ := createTestContext(map[string]string{ + "time[" + op + "]": "2024-01-01T00:00:00Z", + }) + + _, errResp := apirouter.ParseDateFilter(c, "time") + + require.NotNil(t, errResp) + assert.Equal(t, http.StatusBadRequest, errResp.Code) + assert.Contains(t, errResp.Message, "not supported") + }) + } +} + +func TestParseArrayParam(t *testing.T) { + tests := []struct { + name string + fieldName string + queryString string + want []string + }{ + { + name: "no params returns empty slice", + fieldName: "item", + queryString: "", + want: nil, + }, + { + name: "indexed format", + fieldName: "item", + queryString: "item[0]=hello&item[1]=world", + want: []string{"hello", "world"}, + }, + { + name: "indexed format out of order", + fieldName: "item", + queryString: "item[2]=c&item[0]=a&item[1]=b", + want: []string{"a", "b", "c"}, + }, + { + name: "unindexed format", + fieldName: "tag", + queryString: "tag[]=foo&tag[]=bar", + want: []string{"foo", "bar"}, + }, + { + name: "sparse indices", + fieldName: "item", + queryString: "item[0]=first&item[5]=sixth", + want: []string{"first", "sixth"}, + }, + { + name: "ignores non-matching params", + fieldName: "item", + queryString: "item[0]=hello&other=value&item[1]=world", + want: []string{"hello", "world"}, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + c, _ := createTestContextWithQuery(tt.queryString) + + result := apirouter.ParseArrayParam(c, tt.fieldName) + + assert.Equal(t, tt.want, result) + }) + } +} + +func createTestContext(queryParams map[string]string) (*gin.Context, *httptest.ResponseRecorder) { + w := httptest.NewRecorder() + c, _ := gin.CreateTestContext(w) + + values := url.Values{} + for k, v := range queryParams { + if v != "" { + values.Set(k, v) + } + } + + c.Request = httptest.NewRequest(http.MethodGet, "/?"+values.Encode(), nil) + return c, w +} + +func createTestContextWithQuery(queryString string) (*gin.Context, *httptest.ResponseRecorder) { + w := httptest.NewRecorder() + c, _ := gin.CreateTestContext(w) + + path := "/" + if queryString != "" { + path = "/?" + queryString + } + c.Request = httptest.NewRequest(http.MethodGet, path, nil) + return c, w +} + +func timePtr(t time.Time) *time.Time { + return &t +} + +func assertTimePtr(t *testing.T, name string, want, got *time.Time) { + t.Helper() + if want != nil { + require.NotNil(t, got, "%s should not be nil", name) + assert.True(t, want.Equal(*got), "%s mismatch: want %v, got %v", name, want, got) + } else { + assert.Nil(t, got, "%s should be nil", name) + } +} diff --git a/internal/apirouter/tenant_handlers.go b/internal/apirouter/tenant_handlers.go index 0a4a5e492..9fc52983d 100644 --- a/internal/apirouter/tenant_handlers.go +++ b/internal/apirouter/tenant_handlers.go @@ -96,11 +96,24 @@ func (h *TenantHandlers) Retrieve(c *gin.Context) { } func (h *TenantHandlers) List(c *gin.Context) { - // Parse query parameters + // Parse and validate cursors (next/prev are mutually exclusive) + cursors, errResp := ParseCursors(c) + if errResp != nil { + AbortWithError(c, errResp.Code, *errResp) + return + } + + // Parse and validate dir (sort direction) + dir, errResp := ParseDir(c) + if errResp != nil { + AbortWithError(c, errResp.Code, *errResp) + return + } + req := models.ListTenantRequest{ - Next: c.Query("next"), - Prev: c.Query("prev"), - Order: c.Query("order"), + Next: cursors.Next, + Prev: cursors.Prev, + Dir: dir, } // Parse limit if provided diff --git a/internal/logstore/chlogstore/chlogstore.go b/internal/logstore/chlogstore/chlogstore.go index 12cb677a2..1bcdce8ed 100644 --- a/internal/logstore/chlogstore/chlogstore.go +++ b/internal/logstore/chlogstore/chlogstore.go @@ -118,13 +118,21 @@ func buildEventQuery(table string, req driver.ListEventRequest, q pagination.Que args = append(args, req.Topics) } - if req.EventStart != nil { + if req.TimeFilter.GTE != nil { conditions = append(conditions, "event_time >= ?") - args = append(args, *req.EventStart) + args = append(args, *req.TimeFilter.GTE) } - if req.EventEnd != nil { + if req.TimeFilter.LTE != nil { conditions = append(conditions, "event_time <= ?") - args = append(args, *req.EventEnd) + args = append(args, *req.TimeFilter.LTE) + } + if req.TimeFilter.GT != nil { + conditions = append(conditions, "event_time > ?") + args = append(args, *req.TimeFilter.GT) + } + if req.TimeFilter.LT != nil { + conditions = append(conditions, "event_time < ?") + args = append(args, *req.TimeFilter.LT) } if q.CursorPos != "" { @@ -334,13 +342,21 @@ func buildDeliveryQuery(table string, req driver.ListDeliveryEventRequest, q pag args = append(args, req.Topics) } - if req.Start != nil { + if req.TimeFilter.GTE != nil { conditions = append(conditions, "delivery_time >= ?") - args = append(args, *req.Start) + args = append(args, *req.TimeFilter.GTE) } - if req.End != nil { + if req.TimeFilter.LTE != nil { conditions = append(conditions, "delivery_time <= ?") - args = append(args, *req.End) + args = append(args, *req.TimeFilter.LTE) + } + if req.TimeFilter.GT != nil { + conditions = append(conditions, "delivery_time > ?") + args = append(args, *req.TimeFilter.GT) + } + if req.TimeFilter.LT != nil { + conditions = append(conditions, "delivery_time < ?") + args = append(args, *req.TimeFilter.LT) } if q.CursorPos != "" { diff --git a/internal/logstore/driver/driver.go b/internal/logstore/driver/driver.go index dd86fc9d3..fe7cf872d 100644 --- a/internal/logstore/driver/driver.go +++ b/internal/logstore/driver/driver.go @@ -7,6 +7,15 @@ import ( "github.com/hookdeck/outpost/internal/models" ) +// TimeFilter represents time-based filter criteria with support for +// both inclusive (GTE/LTE) and exclusive (GT/LT) comparisons. +type TimeFilter struct { + GTE *time.Time // Greater than or equal (>=) + LTE *time.Time // Less than or equal (<=) + GT *time.Time // Greater than (>) + LT *time.Time // Less than (<) +} + type LogStore interface { ListEvent(context.Context, ListEventRequest) (ListEventResponse, error) ListDeliveryEvent(context.Context, ListDeliveryEventRequest) (ListDeliveryEventResponse, error) @@ -19,8 +28,7 @@ type ListEventRequest struct { Next string Prev string Limit int - EventStart *time.Time // optional - filter events created after this time - EventEnd *time.Time // optional - filter events created before this time + TimeFilter TimeFilter // optional - filter events by time TenantID string // optional - filter by tenant (if empty, returns all tenants) DestinationIDs []string // optional Topics []string // optional @@ -37,8 +45,7 @@ type ListDeliveryEventRequest struct { Next string Prev string Limit int - Start *time.Time // optional - filter deliveries after this time - End *time.Time // optional - filter deliveries before this time + TimeFilter TimeFilter // optional - filter deliveries by time TenantID string // optional - filter by tenant (if empty, returns all tenants) EventID string // optional - filter for specific event DestinationIDs []string // optional diff --git a/internal/logstore/drivertest/crud.go b/internal/logstore/drivertest/crud.go index d4853bc04..45083c556 100644 --- a/internal/logstore/drivertest/crud.go +++ b/internal/logstore/drivertest/crud.go @@ -74,10 +74,10 @@ func testCRUD(t *testing.T, newHarness HarnessMaker) { // Verify via List response, err := logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ - TenantID: tenantID, - EventID: event.ID, - Limit: 10, - Start: &startTime, + TenantID: tenantID, + EventID: event.ID, + Limit: 10, + TimeFilter: driver.TimeFilter{GTE: &startTime}, }) require.NoError(t, err) require.Len(t, response.Data, 1) @@ -138,9 +138,9 @@ func testCRUD(t *testing.T, newHarness HarnessMaker) { // Verify all inserted response, err := logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ - TenantID: tenantID, - Limit: 100, - Start: &startTime, + TenantID: tenantID, + Limit: 100, + TimeFilter: driver.TimeFilter{GTE: &startTime}, }) require.NoError(t, err) // 15 batch + 1 single = 16 @@ -160,7 +160,7 @@ func testCRUD(t *testing.T, newHarness HarnessMaker) { TenantID: tenantID, DestinationIDs: []string{destID}, Limit: 100, - EventStart: &startTime, + TimeFilter: driver.TimeFilter{GTE: &startTime}, }) require.NoError(t, err) require.Len(t, response.Data, len(destinationEvents[destID])) @@ -175,7 +175,7 @@ func testCRUD(t *testing.T, newHarness HarnessMaker) { TenantID: tenantID, DestinationIDs: destIDs, Limit: 100, - EventStart: &startTime, + TimeFilter: driver.TimeFilter{GTE: &startTime}, }) require.NoError(t, err) expectedCount := len(destinationEvents[destIDs[0]]) + len(destinationEvents[destIDs[1]]) @@ -188,7 +188,7 @@ func testCRUD(t *testing.T, newHarness HarnessMaker) { TenantID: tenantID, Topics: []string{topic}, Limit: 100, - EventStart: &startTime, + TimeFilter: driver.TimeFilter{GTE: &startTime}, }) require.NoError(t, err) require.Len(t, response.Data, len(topicEvents[topic])) @@ -199,8 +199,7 @@ func testCRUD(t *testing.T, newHarness HarnessMaker) { eventEnd := baseTime response, err := logStore.ListEvent(ctx, driver.ListEventRequest{ TenantID: tenantID, - EventStart: &eventStart, - EventEnd: &eventEnd, + TimeFilter: driver.TimeFilter{GTE: &eventStart, LTE: &eventEnd}, Limit: 100, }) require.NoError(t, err) @@ -217,7 +216,7 @@ func testCRUD(t *testing.T, newHarness HarnessMaker) { TenantID: tenantID, DestinationIDs: []string{destID}, Limit: 100, - Start: &startTime, + TimeFilter: driver.TimeFilter{GTE: &startTime}, }) require.NoError(t, err) for _, de := range response.Data { @@ -227,10 +226,10 @@ func testCRUD(t *testing.T, newHarness HarnessMaker) { t.Run("ListDeliveryEvent by status", func(t *testing.T) { response, err := logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ - TenantID: tenantID, - Status: "success", - Limit: 100, - Start: &startTime, + TenantID: tenantID, + Status: "success", + Limit: 100, + TimeFilter: driver.TimeFilter{GTE: &startTime}, }) require.NoError(t, err) for _, de := range response.Data { @@ -241,10 +240,10 @@ func testCRUD(t *testing.T, newHarness HarnessMaker) { t.Run("ListDeliveryEvent by topic", func(t *testing.T) { topic := testutil.TestTopics[0] response, err := logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ - TenantID: tenantID, - Topics: []string{topic}, - Limit: 100, - Start: &startTime, + TenantID: tenantID, + Topics: []string{topic}, + Limit: 100, + TimeFilter: driver.TimeFilter{GTE: &startTime}, }) require.NoError(t, err) for _, de := range response.Data { @@ -255,10 +254,10 @@ func testCRUD(t *testing.T, newHarness HarnessMaker) { t.Run("ListDeliveryEvent by event ID", func(t *testing.T) { eventID := "batch_evt_00" response, err := logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ - TenantID: tenantID, - EventID: eventID, - Limit: 100, - Start: &startTime, + TenantID: tenantID, + EventID: eventID, + Limit: 100, + TimeFilter: driver.TimeFilter{GTE: &startTime}, }) require.NoError(t, err) require.Len(t, response.Data, 1) diff --git a/internal/logstore/drivertest/misc.go b/internal/logstore/drivertest/misc.go index 4070ac48e..dc594ef65 100644 --- a/internal/logstore/drivertest/misc.go +++ b/internal/logstore/drivertest/misc.go @@ -86,18 +86,18 @@ func testIsolation(t *testing.T, ctx context.Context, logStore driver.LogStore, t.Run("TenantIsolation", func(t *testing.T) { t.Run("ListDeliveryEvent isolates by tenant", func(t *testing.T) { response, err := logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ - TenantID: tenant1ID, - Limit: 100, - Start: &startTime, + TenantID: tenant1ID, + Limit: 100, + TimeFilter: driver.TimeFilter{GTE: &startTime}, }) require.NoError(t, err) require.Len(t, response.Data, 1) assert.Equal(t, "tenant1-event", response.Data[0].Event.ID) response, err = logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ - TenantID: tenant2ID, - Limit: 100, - Start: &startTime, + TenantID: tenant2ID, + Limit: 100, + TimeFilter: driver.TimeFilter{GTE: &startTime}, }) require.NoError(t, err) require.Len(t, response.Data, 1) @@ -127,7 +127,7 @@ func testIsolation(t *testing.T, ctx context.Context, logStore driver.LogStore, TenantID: "", DestinationIDs: []string{destinationID}, Limit: 100, - EventStart: &startTime, + TimeFilter: driver.TimeFilter{GTE: &startTime}, }) require.NoError(t, err) require.Len(t, response.Data, 2) @@ -145,7 +145,7 @@ func testIsolation(t *testing.T, ctx context.Context, logStore driver.LogStore, TenantID: "", DestinationIDs: []string{destinationID}, Limit: 100, - Start: &startTime, + TimeFilter: driver.TimeFilter{GTE: &startTime}, }) require.NoError(t, err) require.Len(t, response.Data, 2) @@ -222,10 +222,10 @@ func testEdgeCases(t *testing.T, ctx context.Context, logStore driver.LogStore, t.Run("invalid SortOrder uses default (desc)", func(t *testing.T) { response, err := logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ - TenantID: tenantID, - SortOrder: "sideways", - Start: &startTime, - Limit: 10, + TenantID: tenantID, + SortOrder: "sideways", + TimeFilter: driver.TimeFilter{GTE: &startTime}, + Limit: 10, }) require.NoError(t, err) require.Len(t, response.Data, 3) @@ -256,7 +256,7 @@ func testEdgeCases(t *testing.T, ctx context.Context, logStore driver.LogStore, responseNil, err := logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ TenantID: tenantID, DestinationIDs: nil, - Start: &startTime, + TimeFilter: driver.TimeFilter{GTE: &startTime}, Limit: 10, }) require.NoError(t, err) @@ -264,7 +264,7 @@ func testEdgeCases(t *testing.T, ctx context.Context, logStore driver.LogStore, responseEmpty, err := logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ TenantID: tenantID, DestinationIDs: []string{}, - Start: &startTime, + TimeFilter: driver.TimeFilter{GTE: &startTime}, Limit: 10, }) require.NoError(t, err) @@ -312,23 +312,22 @@ func testEdgeCases(t *testing.T, ctx context.Context, logStore driver.LogStore, })) } - t.Run("Start is inclusive (>=)", func(t *testing.T) { + t.Run("GTE is inclusive (>=)", func(t *testing.T) { response, err := logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ - TenantID: tenantID, - Start: &boundaryTime, - Limit: 10, + TenantID: tenantID, + TimeFilter: driver.TimeFilter{GTE: &boundaryTime}, + Limit: 10, }) require.NoError(t, err) require.Len(t, response.Data, 2) }) - t.Run("End is inclusive (<=)", func(t *testing.T) { + t.Run("LTE is inclusive (<=)", func(t *testing.T) { farPast := boundaryTime.Add(-1 * time.Hour) response, err := logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ - TenantID: tenantID, - Start: &farPast, - End: &boundaryTime, - Limit: 10, + TenantID: tenantID, + TimeFilter: driver.TimeFilter{GTE: &farPast, LTE: &boundaryTime}, + Limit: 10, }) require.NoError(t, err) require.Len(t, response.Data, 2) @@ -354,9 +353,9 @@ func testEdgeCases(t *testing.T, ctx context.Context, logStore driver.LogStore, t.Run("modifying ListDeliveryEvent result doesn't affect subsequent queries", func(t *testing.T) { response1, err := logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ - TenantID: tenantID, - Limit: 10, - Start: &startTime, + TenantID: tenantID, + Limit: 10, + TimeFilter: driver.TimeFilter{GTE: &startTime}, }) require.NoError(t, err) require.Len(t, response1.Data, 1) @@ -365,9 +364,9 @@ func testEdgeCases(t *testing.T, ctx context.Context, logStore driver.LogStore, response1.Data[0].Event.ID = "MODIFIED" response2, err := logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ - TenantID: tenantID, - Limit: 10, - Start: &startTime, + TenantID: tenantID, + Limit: 10, + TimeFilter: driver.TimeFilter{GTE: &startTime}, }) require.NoError(t, err) require.Len(t, response2.Data, 1) @@ -416,9 +415,9 @@ func testEdgeCases(t *testing.T, ctx context.Context, logStore driver.LogStore, // Assert: still exactly 1 record response, err := logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ - TenantID: tenantID, - Limit: 100, - Start: &startTime, + TenantID: tenantID, + Limit: 100, + TimeFilter: driver.TimeFilter{GTE: &startTime}, }) require.NoError(t, err) require.Len(t, response.Data, 1, "concurrent duplicate inserts should result in exactly 1 record") @@ -441,11 +440,11 @@ func testCursorValidation(t *testing.T, ctx context.Context, logStore driver.Log for _, tc := range testCases { t.Run(tc.name, func(t *testing.T) { _, err := logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ - TenantID: tenantID, - SortOrder: "desc", - Next: tc.cursor, - Start: &startTime, - Limit: 10, + TenantID: tenantID, + SortOrder: "desc", + Next: tc.cursor, + TimeFilter: driver.TimeFilter{GTE: &startTime}, + Limit: 10, }) require.Error(t, err) assert.True(t, errors.Is(err, cursor.ErrInvalidCursor), "expected cursor.ErrInvalidCursor, got: %v", err) @@ -480,20 +479,20 @@ func testCursorValidation(t *testing.T, ctx context.Context, logStore driver.Log t.Run("delivery_time desc", func(t *testing.T) { page1, err := logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ - TenantID: tenantID, - SortOrder: "desc", - Start: &startTime, - Limit: 2, + TenantID: tenantID, + SortOrder: "desc", + TimeFilter: driver.TimeFilter{GTE: &startTime}, + Limit: 2, }) require.NoError(t, err) require.NotEmpty(t, page1.Next) page2, err := logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ - TenantID: tenantID, - SortOrder: "desc", - Next: page1.Next, - Start: &startTime, - Limit: 2, + TenantID: tenantID, + SortOrder: "desc", + Next: page1.Next, + TimeFilter: driver.TimeFilter{GTE: &startTime}, + Limit: 2, }) require.NoError(t, err) require.NotEmpty(t, page2.Data) @@ -501,20 +500,20 @@ func testCursorValidation(t *testing.T, ctx context.Context, logStore driver.Log t.Run("delivery_time asc", func(t *testing.T) { page1, err := logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ - TenantID: tenantID, - SortOrder: "asc", - Start: &startTime, - Limit: 2, + TenantID: tenantID, + SortOrder: "asc", + TimeFilter: driver.TimeFilter{GTE: &startTime}, + Limit: 2, }) require.NoError(t, err) require.NotEmpty(t, page1.Next) page2, err := logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ - TenantID: tenantID, - SortOrder: "asc", - Next: page1.Next, - Start: &startTime, - Limit: 2, + TenantID: tenantID, + SortOrder: "asc", + Next: page1.Next, + TimeFilter: driver.TimeFilter{GTE: &startTime}, + Limit: 2, }) require.NoError(t, err) require.NotEmpty(t, page2.Data) diff --git a/internal/logstore/drivertest/pagination.go b/internal/logstore/drivertest/pagination.go index 0ba297224..6a48877cd 100644 --- a/internal/logstore/drivertest/pagination.go +++ b/internal/logstore/drivertest/pagination.go @@ -79,12 +79,12 @@ func testPagination(t *testing.T, newHarness HarnessMaker) { List: func(ctx context.Context, opts paginationtest.ListOpts) (paginationtest.ListResult[*models.DeliveryEvent], error) { res, err := logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ - TenantID: tenantID, - Limit: opts.Limit, - SortOrder: opts.Order, - Next: opts.Next, - Prev: opts.Prev, - Start: &farPast, + TenantID: tenantID, + Limit: opts.Limit, + SortOrder: opts.Order, + Next: opts.Next, + Prev: opts.Prev, + TimeFilter: driver.TimeFilter{GTE: &farPast}, }) if err != nil { return paginationtest.ListResult[*models.DeliveryEvent]{}, err @@ -171,7 +171,7 @@ func testPagination(t *testing.T, newHarness HarnessMaker) { SortOrder: opts.Order, Next: opts.Next, Prev: opts.Prev, - Start: &farPast, + TimeFilter: driver.TimeFilter{GTE: &farPast}, }) if err != nil { return paginationtest.ListResult[*models.DeliveryEvent]{}, err @@ -255,7 +255,7 @@ func testPagination(t *testing.T, newHarness HarnessMaker) { SortOrder: opts.Order, Next: opts.Next, Prev: opts.Prev, - EventStart: &farPast, + TimeFilter: driver.TimeFilter{GTE: &farPast}, }) if err != nil { return paginationtest.ListResult[*models.Event]{}, err @@ -342,7 +342,7 @@ func testPagination(t *testing.T, newHarness HarnessMaker) { SortOrder: opts.Order, Next: opts.Next, Prev: opts.Prev, - EventStart: &farPast, + TimeFilter: driver.TimeFilter{GTE: &farPast}, }) if err != nil { return paginationtest.ListResult[*models.Event]{}, err @@ -369,4 +369,416 @@ func testPagination(t *testing.T, newHarness HarnessMaker) { suite.Run(t) }) + + // Test cursor pagination combined with time filters. + // These tests verify that cursors work correctly when used alongside + // time-based filters (GTE, LTE, GT, LT), which is critical for + // "paginate within a time window" use cases. + // + // IMPORTANT: ListDeliveryEvent filters by DELIVERY time, ListEvent filters by EVENT time. + // In this test, delivery_time = event_time + 100ms. + t.Run("TimeFilterWithCursor", func(t *testing.T) { + tenantID := idgen.String() + destinationID := idgen.Destination() + idPrefix := idgen.String()[:8] + + // Create 20 events with times spread across different ranges: + // - Events 0-4: far past (should be excluded by GTE filter) + // - Events 5-14: within time window (should be included) + // - Events 15-19: far future (should be excluded by LTE filter) + // + // Event times are spaced 2 minutes apart within the window. + // Delivery times are 1 second after event times (not sub-second) + // to ensure GT/LT tests work consistently across databases. + eventWindowStart := baseTime.Add(-10 * time.Minute) + eventWindowEnd := baseTime.Add(10 * time.Minute) + // Delivery window accounts for the 1 second offset + deliveryWindowStart := eventWindowStart.Add(time.Second) + deliveryWindowEnd := eventWindowEnd.Add(time.Second) + + var allEvents []*models.DeliveryEvent + for i := range 20 { + var eventTime time.Time + switch { + case i < 5: + // Far past: outside window (before eventWindowStart) + eventTime = eventWindowStart.Add(-time.Duration(5-i) * time.Hour) + case i < 15: + // Within window: eventWindowStart to eventWindowEnd + offset := time.Duration(i-5) * 2 * time.Minute + eventTime = eventWindowStart.Add(offset) + default: + // Far future: outside window (after eventWindowEnd) + eventTime = eventWindowEnd.Add(time.Duration(i-14) * time.Hour) + } + + deliveryTime := eventTime.Add(time.Second) + + event := &models.Event{ + ID: fmt.Sprintf("%s_evt_%03d", idPrefix, i), + TenantID: tenantID, + DestinationID: destinationID, + Topic: "test.topic", + EligibleForRetry: true, + Time: eventTime, + Metadata: map[string]string{}, + Data: map[string]any{}, + } + delivery := &models.Delivery{ + ID: fmt.Sprintf("%s_del_%03d", idPrefix, i), + EventID: event.ID, + DestinationID: destinationID, + Status: "success", + Time: deliveryTime, + Code: "200", + } + allEvents = append(allEvents, &models.DeliveryEvent{ + ID: fmt.Sprintf("%s_de_%03d", idPrefix, i), + DestinationID: destinationID, + Event: *event, + Delivery: delivery, + }) + } + + require.NoError(t, logStore.InsertManyDeliveryEvent(ctx, allEvents)) + require.NoError(t, h.FlushWrites(ctx)) + + t.Run("paginate within time-bounded window", func(t *testing.T) { + // Paginate through deliveries within the window with limit=3 + // ListDeliveryEvent filters by DELIVERY time, not event time. + // Should only see deliveries 5-14 (10 total), not 0-4 or 15-19 + var collectedIDs []string + var nextCursor string + pageCount := 0 + + for { + res, err := logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ + TenantID: tenantID, + Limit: 3, + SortOrder: "asc", + Next: nextCursor, + TimeFilter: driver.TimeFilter{GTE: &deliveryWindowStart, LTE: &deliveryWindowEnd}, + }) + require.NoError(t, err) + + for _, de := range res.Data { + collectedIDs = append(collectedIDs, de.Event.ID) + } + + pageCount++ + if res.Next == "" { + break + } + nextCursor = res.Next + + // Safety: prevent infinite loop + if pageCount > 10 { + t.Fatal("too many pages") + } + } + + // Should have collected exactly deliveries 5-14 + require.Len(t, collectedIDs, 10, "should have 10 deliveries in window") + for i, id := range collectedIDs { + expectedID := fmt.Sprintf("%s_evt_%03d", idPrefix, i+5) + require.Equal(t, expectedID, id, "delivery %d mismatch", i) + } + require.Equal(t, 4, pageCount, "should take 4 pages (3+3+3+1)") + }) + + t.Run("cursor excludes deliveries outside time filter", func(t *testing.T) { + // First page with no time filter gets all deliveries + resAll, err := logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ + TenantID: tenantID, + Limit: 5, + SortOrder: "asc", + TimeFilter: driver.TimeFilter{GTE: &farPast}, + }) + require.NoError(t, err) + require.Len(t, resAll.Data, 5) + + // Use the cursor but add a time filter that excludes some results + // The cursor points to position after delivery 4 (far past deliveries) + // But with deliveryWindowStart filter, we should start from delivery 5 + res, err := logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ + TenantID: tenantID, + Limit: 5, + SortOrder: "asc", + Next: resAll.Next, + TimeFilter: driver.TimeFilter{GTE: &deliveryWindowStart, LTE: &deliveryWindowEnd}, + }) + require.NoError(t, err) + + // Results should respect the time filter (on delivery time) + for _, de := range res.Data { + require.True(t, !de.Delivery.Time.Before(deliveryWindowStart), "delivery time should be >= deliveryWindowStart") + require.True(t, !de.Delivery.Time.After(deliveryWindowEnd), "delivery time should be <= deliveryWindowEnd") + } + }) + + t.Run("delivery time filter with GT/LT operators", func(t *testing.T) { + // Test exclusive bounds (GT/LT instead of GTE/LTE) on delivery time + // Use delivery times slightly after delivery 5 and slightly before delivery 14 + gtTime := allEvents[5].Delivery.Time.Add(time.Second) // After delivery 5, before delivery 6 + ltTime := allEvents[14].Delivery.Time.Add(-time.Second) // Before delivery 14, after delivery 13 + + var collectedIDs []string + var nextCursor string + + for { + res, err := logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ + TenantID: tenantID, + Limit: 3, + SortOrder: "asc", + Next: nextCursor, + TimeFilter: driver.TimeFilter{GT: >Time, LT: <Time}, + }) + require.NoError(t, err) + + for _, de := range res.Data { + collectedIDs = append(collectedIDs, de.Event.ID) + } + + if res.Next == "" { + break + } + nextCursor = res.Next + } + + // Should have events 6-13 (8 events) + require.Len(t, collectedIDs, 8, "should have 8 events in GT/LT range") + for i, id := range collectedIDs { + expectedID := fmt.Sprintf("%s_evt_%03d", idPrefix, i+6) + require.Equal(t, expectedID, id, "event %d mismatch", i) + } + }) + + t.Run("GT/LT exclude exact timestamp", func(t *testing.T) { + // Verify that GT excludes the exact timestamp (not >=) + // and LT excludes the exact timestamp (not <=). + // + // We truncate times to second precision to ensure consistent + // comparison across databases with different timestamp precision + // (PostgreSQL microseconds, ClickHouse DateTime64, etc.). + // + // Important: ListDeliveryEvent filters by DELIVERY time, not event time. + + // First, retrieve all deliveries to find delivery 10's time + res, err := logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ + TenantID: tenantID, + Limit: 100, + SortOrder: "asc", + TimeFilter: driver.TimeFilter{ + GTE: &farPast, + }, + }) + require.NoError(t, err) + require.GreaterOrEqual(t, len(res.Data), 11, "need at least 11 deliveries") + + // Find delivery 10's stored delivery time, truncated to seconds + var storedDelivery10Time time.Time + for _, de := range res.Data { + if de.Event.ID == allEvents[10].Event.ID { + storedDelivery10Time = de.Delivery.Time.Truncate(time.Second) + break + } + } + require.False(t, storedDelivery10Time.IsZero(), "should find delivery 10") + + // GT with exact time should exclude delivery 10 + resGT, err := logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ + TenantID: tenantID, + Limit: 100, + SortOrder: "asc", + TimeFilter: driver.TimeFilter{GT: &storedDelivery10Time}, + }) + require.NoError(t, err) + + for _, de := range resGT.Data { + deTimeTrunc := de.Delivery.Time.Truncate(time.Second) + require.True(t, deTimeTrunc.After(storedDelivery10Time), + "GT filter should exclude delivery with exact timestamp, got delivery %s with time %v (filter time: %v)", + de.Delivery.ID, deTimeTrunc, storedDelivery10Time) + } + + // LT with exact time should exclude delivery 10 + resLT, err := logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ + TenantID: tenantID, + Limit: 100, + SortOrder: "asc", + TimeFilter: driver.TimeFilter{LT: &storedDelivery10Time}, + }) + require.NoError(t, err) + + for _, de := range resLT.Data { + deTimeTrunc := de.Delivery.Time.Truncate(time.Second) + require.True(t, deTimeTrunc.Before(storedDelivery10Time), + "LT filter should exclude delivery with exact timestamp, got delivery %s with time %v (filter time: %v)", + de.Delivery.ID, deTimeTrunc, storedDelivery10Time) + } + + // Verify delivery 10 is included with GTE/LTE (inclusive bounds) + resGTE, err := logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ + TenantID: tenantID, + Limit: 100, + SortOrder: "asc", + TimeFilter: driver.TimeFilter{GTE: &storedDelivery10Time, LTE: &storedDelivery10Time}, + }) + require.NoError(t, err) + require.GreaterOrEqual(t, len(resGTE.Data), 1, "GTE/LTE with same time should include delivery at that second") + }) + + t.Run("prev cursor respects time filter", func(t *testing.T) { + // Get first page (ListDeliveryEvent filters by delivery time) + res1, err := logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ + TenantID: tenantID, + Limit: 3, + SortOrder: "asc", + TimeFilter: driver.TimeFilter{GTE: &deliveryWindowStart, LTE: &deliveryWindowEnd}, + }) + require.NoError(t, err) + require.NotEmpty(t, res1.Next) + + // Get second page + res2, err := logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ + TenantID: tenantID, + Limit: 3, + SortOrder: "asc", + Next: res1.Next, + TimeFilter: driver.TimeFilter{GTE: &deliveryWindowStart, LTE: &deliveryWindowEnd}, + }) + require.NoError(t, err) + require.NotEmpty(t, res2.Prev) + + // Go back to first page using prev cursor + resPrev, err := logStore.ListDeliveryEvent(ctx, driver.ListDeliveryEventRequest{ + TenantID: tenantID, + Limit: 3, + SortOrder: "asc", + Prev: res2.Prev, + TimeFilter: driver.TimeFilter{GTE: &deliveryWindowStart, LTE: &deliveryWindowEnd}, + }) + require.NoError(t, err) + + // Should get same results as first page + require.Len(t, resPrev.Data, len(res1.Data)) + for i := range res1.Data { + require.Equal(t, res1.Data[i].Event.ID, resPrev.Data[i].Event.ID) + } + }) + + t.Run("ListEvent with time filter pagination", func(t *testing.T) { + // Same test pattern for ListEvent + var collectedIDs []string + var nextCursor string + pageCount := 0 + + for { + res, err := logStore.ListEvent(ctx, driver.ListEventRequest{ + TenantID: tenantID, + Limit: 3, + SortOrder: "asc", + Next: nextCursor, + TimeFilter: driver.TimeFilter{GTE: &eventWindowStart, LTE: &eventWindowEnd}, + }) + require.NoError(t, err) + + for _, e := range res.Data { + collectedIDs = append(collectedIDs, e.ID) + } + + pageCount++ + if res.Next == "" { + break + } + nextCursor = res.Next + + if pageCount > 10 { + t.Fatal("too many pages") + } + } + + // Should have collected exactly events 5-14 + require.Len(t, collectedIDs, 10, "should have 10 events in window") + for i, id := range collectedIDs { + expectedID := fmt.Sprintf("%s_evt_%03d", idPrefix, i+5) + require.Equal(t, expectedID, id, "event %d mismatch", i) + } + require.Equal(t, 4, pageCount, "should take 4 pages (3+3+3+1)") + }) + + t.Run("ListEvent GT/LT exclude exact timestamp", func(t *testing.T) { + // Verify that GT excludes the exact timestamp (not >=) + // and LT excludes the exact timestamp (not <=). + // + // We truncate times to second precision to ensure consistent + // comparison across databases with different timestamp precision. + // + // ListEvent filters by EVENT time. + + // First, retrieve event 10's stored time from the database + res, err := logStore.ListEvent(ctx, driver.ListEventRequest{ + TenantID: tenantID, + Limit: 100, + SortOrder: "asc", + TimeFilter: driver.TimeFilter{ + GTE: &farPast, + }, + }) + require.NoError(t, err) + require.GreaterOrEqual(t, len(res.Data), 11, "need at least 11 events") + + // Find event 10's stored event time, truncated to seconds + var storedEvent10Time time.Time + for _, e := range res.Data { + if e.ID == allEvents[10].Event.ID { + storedEvent10Time = e.Time.Truncate(time.Second) + break + } + } + require.False(t, storedEvent10Time.IsZero(), "should find event 10") + + // GT with exact time should exclude event 10 + resGT, err := logStore.ListEvent(ctx, driver.ListEventRequest{ + TenantID: tenantID, + Limit: 100, + SortOrder: "asc", + TimeFilter: driver.TimeFilter{GT: &storedEvent10Time}, + }) + require.NoError(t, err) + + for _, e := range resGT.Data { + eTimeTrunc := e.Time.Truncate(time.Second) + require.True(t, eTimeTrunc.After(storedEvent10Time), + "GT filter should exclude event with exact timestamp, got event %s with time %v (filter time: %v)", + e.ID, eTimeTrunc, storedEvent10Time) + } + + // LT with exact time should exclude event 10 + resLT, err := logStore.ListEvent(ctx, driver.ListEventRequest{ + TenantID: tenantID, + Limit: 100, + SortOrder: "asc", + TimeFilter: driver.TimeFilter{LT: &storedEvent10Time}, + }) + require.NoError(t, err) + + for _, e := range resLT.Data { + eTimeTrunc := e.Time.Truncate(time.Second) + require.True(t, eTimeTrunc.Before(storedEvent10Time), + "LT filter should exclude event with exact timestamp, got event %s with time %v (filter time: %v)", + e.ID, eTimeTrunc, storedEvent10Time) + } + + // Verify event 10 is included with GTE/LTE (inclusive bounds) + resGTE, err := logStore.ListEvent(ctx, driver.ListEventRequest{ + TenantID: tenantID, + Limit: 100, + SortOrder: "asc", + TimeFilter: driver.TimeFilter{GTE: &storedEvent10Time, LTE: &storedEvent10Time}, + }) + require.NoError(t, err) + require.GreaterOrEqual(t, len(resGTE.Data), 1, "GTE/LTE with same time should include event at that second") + }) + }) } diff --git a/internal/logstore/logstore.go b/internal/logstore/logstore.go index 5ac6c7c8e..84a45314a 100644 --- a/internal/logstore/logstore.go +++ b/internal/logstore/logstore.go @@ -13,6 +13,7 @@ import ( "github.com/jackc/pgx/v5/pgxpool" ) +type TimeFilter = driver.TimeFilter type ListEventRequest = driver.ListEventRequest type ListEventResponse = driver.ListEventResponse type ListDeliveryEventRequest = driver.ListDeliveryEventRequest diff --git a/internal/logstore/memlogstore/memlogstore.go b/internal/logstore/memlogstore/memlogstore.go index a68abcda6..fc6c8014b 100644 --- a/internal/logstore/memlogstore/memlogstore.go +++ b/internal/logstore/memlogstore/memlogstore.go @@ -175,10 +175,16 @@ func (s *memLogStore) matchesEventFilter(event *models.Event, req driver.ListEve } } - if req.EventStart != nil && event.Time.Before(*req.EventStart) { + if req.TimeFilter.GTE != nil && event.Time.Before(*req.TimeFilter.GTE) { return false } - if req.EventEnd != nil && event.Time.After(*req.EventEnd) { + if req.TimeFilter.LTE != nil && event.Time.After(*req.TimeFilter.LTE) { + return false + } + if req.TimeFilter.GT != nil && !event.Time.After(*req.TimeFilter.GT) { + return false + } + if req.TimeFilter.LT != nil && !event.Time.Before(*req.TimeFilter.LT) { return false } @@ -391,10 +397,16 @@ func (s *memLogStore) matchesFilter(de *models.DeliveryEvent, req driver.ListDel } } - if req.Start != nil && de.Delivery.Time.Before(*req.Start) { + if req.TimeFilter.GTE != nil && de.Delivery.Time.Before(*req.TimeFilter.GTE) { + return false + } + if req.TimeFilter.LTE != nil && de.Delivery.Time.After(*req.TimeFilter.LTE) { + return false + } + if req.TimeFilter.GT != nil && !de.Delivery.Time.After(*req.TimeFilter.GT) { return false } - if req.End != nil && de.Delivery.Time.After(*req.End) { + if req.TimeFilter.LT != nil && !de.Delivery.Time.Before(*req.TimeFilter.LT) { return false } diff --git a/internal/logstore/pglogstore/pglogstore.go b/internal/logstore/pglogstore/pglogstore.go index 9e026f3f3..90f7319f6 100644 --- a/internal/logstore/pglogstore/pglogstore.go +++ b/internal/logstore/pglogstore/pglogstore.go @@ -88,7 +88,7 @@ func (s *logStore) ListEvent(ctx context.Context, req driver.ListEventRequest) ( } func buildEventQuery(req driver.ListEventRequest, q pagination.QueryInput) (string, []any) { - cursorCondition := fmt.Sprintf("AND ($6::text = '' OR time_id %s $6::text)", q.Compare) + cursorCondition := fmt.Sprintf("AND ($8::text = '' OR time_id %s $8::text)", q.Compare) orderByClause := fmt.Sprintf("time %s, id %s", strings.ToUpper(q.SortDir), strings.ToUpper(q.SortDir)) query := fmt.Sprintf(` @@ -108,19 +108,23 @@ func buildEventQuery(req driver.ListEventRequest, q pagination.QueryInput) (stri AND (array_length($3::text[], 1) IS NULL OR topic = ANY($3)) AND ($4::timestamptz IS NULL OR time >= $4) AND ($5::timestamptz IS NULL OR time <= $5) + AND ($6::timestamptz IS NULL OR time > $6) + AND ($7::timestamptz IS NULL OR time < $7) %s ORDER BY %s - LIMIT $7 + LIMIT $9 `, cursorCondition, orderByClause) args := []any{ req.TenantID, // $1 req.DestinationIDs, // $2 req.Topics, // $3 - req.EventStart, // $4 - req.EventEnd, // $5 - q.CursorPos, // $6 - q.Limit, // $7 + req.TimeFilter.GTE, // $4 + req.TimeFilter.LTE, // $5 + req.TimeFilter.GT, // $6 + req.TimeFilter.LT, // $7 + q.CursorPos, // $8 + q.Limit, // $9 } return query, args @@ -235,7 +239,7 @@ func (s *logStore) ListDeliveryEvent(ctx context.Context, req driver.ListDeliver } func buildDeliveryQuery(req driver.ListDeliveryEventRequest, q pagination.QueryInput) (string, []any) { - cursorCondition := fmt.Sprintf("AND ($8::text = '' OR idx.time_delivery_id %s $8::text)", q.Compare) + cursorCondition := fmt.Sprintf("AND ($10::text = '' OR idx.time_delivery_id %s $10::text)", q.Compare) orderByClause := fmt.Sprintf("idx.delivery_time %s, idx.delivery_id %s", strings.ToUpper(q.SortDir), strings.ToUpper(q.SortDir)) query := fmt.Sprintf(` @@ -266,9 +270,11 @@ func buildDeliveryQuery(req driver.ListDeliveryEventRequest, q pagination.QueryI AND (array_length($5::text[], 1) IS NULL OR idx.topic = ANY($5)) AND ($6::timestamptz IS NULL OR idx.delivery_time >= $6) AND ($7::timestamptz IS NULL OR idx.delivery_time <= $7) + AND ($8::timestamptz IS NULL OR idx.delivery_time > $8) + AND ($9::timestamptz IS NULL OR idx.delivery_time < $9) %s ORDER BY %s - LIMIT $9 + LIMIT $11 `, cursorCondition, orderByClause) args := []any{ @@ -277,10 +283,12 @@ func buildDeliveryQuery(req driver.ListDeliveryEventRequest, q pagination.QueryI req.DestinationIDs, // $3 req.Status, // $4 req.Topics, // $5 - req.Start, // $6 - req.End, // $7 - q.CursorPos, // $8 - q.Limit, // $9 + req.TimeFilter.GTE, // $6 + req.TimeFilter.LTE, // $7 + req.TimeFilter.GT, // $8 + req.TimeFilter.LT, // $9 + q.CursorPos, // $10 + q.Limit, // $11 } return query, args diff --git a/internal/models/entity.go b/internal/models/entity.go index ecf776ed8..d60547bd1 100644 --- a/internal/models/entity.go +++ b/internal/models/entity.go @@ -21,7 +21,7 @@ type EntityStore interface { RetrieveTenant(ctx context.Context, tenantID string) (*Tenant, error) UpsertTenant(ctx context.Context, tenant Tenant) error DeleteTenant(ctx context.Context, tenantID string) error - ListTenant(ctx context.Context, req ListTenantRequest) (*ListTenantResponse, error) + ListTenant(ctx context.Context, req ListTenantRequest) (*TenantPaginatedResult, error) ListDestinationByTenant(ctx context.Context, tenantID string, options ...ListDestinationByTenantOpts) ([]Destination, error) RetrieveDestination(ctx context.Context, tenantID, destinationID string) (*Destination, error) CreateDestination(ctx context.Context, destination Destination) error @@ -48,15 +48,23 @@ type ListTenantRequest struct { Limit int // Number of results per page (default: 20) Next string // Cursor for next page Prev string // Cursor for previous page - Order string // Sort order: "asc" or "desc" (default: "desc") + Dir string // Sort direction: "asc" or "desc" (default: "desc") } -// ListTenantResponse contains the paginated list of tenants. -type ListTenantResponse struct { - Data []Tenant `json:"data"` - Next string `json:"next"` - Prev string `json:"prev"` - Count int `json:"count"` +// SeekPagination represents cursor-based pagination metadata for list responses. +type SeekPagination struct { + OrderBy string `json:"order_by"` + Dir string `json:"dir"` + Limit int `json:"limit"` + Next *string `json:"next"` + Prev *string `json:"prev"` +} + +// TenantPaginatedResult contains the paginated list of tenants. +type TenantPaginatedResult struct { + Models []Tenant `json:"models"` + Pagination SeekPagination `json:"pagination"` + Count int `json:"count"` } type entityStoreImpl struct { @@ -331,7 +339,7 @@ const ( ) // ListTenant returns a paginated list of tenants using RediSearch. -func (s *entityStoreImpl) ListTenant(ctx context.Context, req ListTenantRequest) (*ListTenantResponse, error) { +func (s *entityStoreImpl) ListTenant(ctx context.Context, req ListTenantRequest) (*TenantPaginatedResult, error) { if !s.listTenantSupported { return nil, ErrListTenantNotSupported } @@ -350,12 +358,12 @@ func (s *entityStoreImpl) ListTenant(ctx context.Context, req ListTenantRequest) limit = maxListTenantLimit } - // Validate and apply order - order := req.Order - if order == "" { - order = "desc" + // Validate and apply dir (sort direction) + dir := req.Dir + if dir == "" { + dir = "desc" } - if order != "asc" && order != "desc" { + if dir != "asc" && dir != "desc" { return nil, ErrInvalidOrder } @@ -367,7 +375,7 @@ func (s *entityStoreImpl) ListTenant(ctx context.Context, req ListTenantRequest) // Use pagination package for cursor-based pagination with n+1 pattern result, err := pagination.Run(ctx, pagination.Config[Tenant]{ Limit: limit, - Order: order, + Order: dir, Next: req.Next, Prev: req.Prev, Cursor: pagination.Cursor[Tenant]{ @@ -424,11 +432,25 @@ func (s *entityStoreImpl) ListTenant(ctx context.Context, req ListTenantRequest) _, totalCount, _ = s.parseSearchResult(ctx, countResult) } - return &ListTenantResponse{ - Data: tenants, + // Convert empty cursors to nil pointers (Hookdeck returns null for empty cursors) + var nextCursor, prevCursor *string + if result.Next != "" { + nextCursor = &result.Next + } + if result.Prev != "" { + prevCursor = &result.Prev + } + + return &TenantPaginatedResult{ + Models: tenants, + Pagination: SeekPagination{ + OrderBy: "created_at", + Dir: dir, + Limit: limit, + Next: nextCursor, + Prev: prevCursor, + }, Count: totalCount, - Next: result.Next, - Prev: result.Prev, }, nil } @@ -873,7 +895,7 @@ func (s *entityStoreImpl) parseTenantTopics(destinationSummaryList []Destination } if all { - return s.availableTopics + return []string{"*"} } topics := make([]string, 0, len(topicsSet)) diff --git a/internal/models/entity_test.go b/internal/models/entity_test.go index 4cb489948..d6c8f3ab8 100644 --- a/internal/models/entity_test.go +++ b/internal/models/entity_test.go @@ -243,17 +243,25 @@ func runListTenantPaginationSuite(t *testing.T, factory RedisClientFactory, depl List: func(ctx context.Context, opts paginationtest.ListOpts) (paginationtest.ListResult[models.Tenant], error) { resp, err := entityStore.ListTenant(ctx, models.ListTenantRequest{ Limit: opts.Limit, - Order: opts.Order, + Dir: opts.Order, Next: opts.Next, Prev: opts.Prev, }) if err != nil { return paginationtest.ListResult[models.Tenant]{}, err } + // Convert *string cursors to string for test framework + var next, prev string + if resp.Pagination.Next != nil { + next = *resp.Pagination.Next + } + if resp.Pagination.Prev != nil { + prev = *resp.Pagination.Prev + } return paginationtest.ListResult[models.Tenant]{ - Items: resp.Data, - Next: resp.Next, - Prev: resp.Prev, + Items: resp.Models, + Next: next, + Prev: prev, }, nil }, diff --git a/internal/models/entitysuite_test.go b/internal/models/entitysuite_test.go index 2378c783e..ff025a3df 100644 --- a/internal/models/entitysuite_test.go +++ b/internal/models/entitysuite_test.go @@ -542,10 +542,12 @@ func (s *EntityTestSuite) TestMultiDestinationRetrieveTenantDestinationsCount() func (s *EntityTestSuite) TestMultiDestinationRetrieveTenantTopics() { data := s.setupMultiDestination() + // destinations[0] has topics ["*"], so tenant.Topics should be ["*"] tenant, err := s.entityStore.RetrieveTenant(s.ctx, data.tenant.ID) require.NoError(s.T(), err) - require.Equal(s.T(), []string{"user.created", "user.deleted", "user.updated"}, tenant.Topics) + require.Equal(s.T(), []string{"*"}, tenant.Topics) + // After deleting the wildcard destination, tenant.Topics should aggregate remaining topics require.NoError(s.T(), s.entityStore.DeleteDestination(s.ctx, data.tenant.ID, data.destinations[0].ID)) tenant, err = s.entityStore.RetrieveTenant(s.ctx, data.tenant.ID) require.NoError(s.T(), err) @@ -1181,9 +1183,9 @@ func (s *ListTenantTestSuite) SetupSuite() { // Test empty list BEFORE creating any data resp, err := s.entityStore.ListTenant(s.ctx, models.ListTenantRequest{}) require.NoError(s.T(), err) - require.Empty(s.T(), resp.Data, "should be empty before creating data") - require.Empty(s.T(), resp.Next) - require.Empty(s.T(), resp.Prev) + require.Empty(s.T(), resp.Models, "should be empty before creating data") + require.Nil(s.T(), resp.Pagination.Next) + require.Nil(s.T(), resp.Pagination.Prev) // Create 25 tenants for pagination tests s.tenants = make([]models.Tenant, 25) @@ -1218,10 +1220,14 @@ func (s *ListTenantTestSuite) TestListTenantEnrichment() { resp1, err := s.entityStore.ListTenant(s.ctx, models.ListTenantRequest{Limit: 2}) require.NoError(t, err) assert.Equal(t, 25, resp1.Count, "count should be total tenants, not page size") - assert.Len(t, resp1.Data, 2, "data should respect limit") + assert.Len(t, resp1.Models, 2, "data should respect limit") // Second page - count should still be total - resp2, err := s.entityStore.ListTenant(s.ctx, models.ListTenantRequest{Limit: 2, Next: resp1.Next}) + var nextCursor string + if resp1.Pagination.Next != nil { + nextCursor = *resp1.Pagination.Next + } + resp2, err := s.entityStore.ListTenant(s.ctx, models.ListTenantRequest{Limit: 2, Next: nextCursor}) require.NoError(t, err) assert.Equal(t, 25, resp2.Count, "count should remain total across pages") }) @@ -1231,7 +1237,7 @@ func (s *ListTenantTestSuite) TestListTenantEnrichment() { require.NoError(t, err) assert.Equal(t, 25, resp.Count, "count should be tenants only") - for _, tenant := range resp.Data { + for _, tenant := range resp.Models { assert.NotContains(t, tenant.ID, "dest_", "destination should not appear in tenant list") } }) @@ -1242,9 +1248,9 @@ func (s *ListTenantTestSuite) TestListTenantEnrichment() { // tenantWithDests has 2 destinations from SetupSuite var tenantWithDests *models.Tenant - for i := range resp.Data { - if resp.Data[i].ID == s.tenantWithDests.ID { - tenantWithDests = &resp.Data[i] + for i := range resp.Models { + if resp.Models[i].ID == s.tenantWithDests.ID { + tenantWithDests = &resp.Models[i] break } } @@ -1254,9 +1260,9 @@ func (s *ListTenantTestSuite) TestListTenantEnrichment() { // Verify tenants without destinations have 0 count var tenantWithoutDests *models.Tenant - for i := range resp.Data { - if resp.Data[i].ID != s.tenantWithDests.ID { - tenantWithoutDests = &resp.Data[i] + for i := range resp.Models { + if resp.Models[i].ID != s.tenantWithDests.ID { + tenantWithoutDests = &resp.Models[i] break } } @@ -1293,7 +1299,7 @@ func (s *ListTenantTestSuite) TestListTenantExcludesDeleted() { assert.Equal(t, initialCount+1, resp.Count) // Verify deleted tenant is not in results - for _, tenant := range resp.Data { + for _, tenant := range resp.Models { assert.NotEqual(t, tenant1.ID, tenant.ID, "deleted tenant should not appear") } @@ -1328,15 +1334,15 @@ func (s *ListTenantTestSuite) TestListTenantExcludesDeleted() { // NOT the deleted ones (index 3,4) which would have been first if not filtered resp, err := s.entityStore.ListTenant(s.ctx, models.ListTenantRequest{ Limit: 2, - Order: "desc", + Dir: "desc", }) require.NoError(t, err) - require.GreaterOrEqual(t, len(resp.Data), 2, "should get at least 2 tenants") + require.GreaterOrEqual(t, len(resp.Models), 2, "should get at least 2 tenants") // Verify deleted tenants don't appear in the first 2 results - for i := 0; i < 2 && i < len(resp.Data); i++ { - assert.NotEqual(t, tenantIDs[3], resp.Data[i].ID, "deleted tenant should not appear") - assert.NotEqual(t, tenantIDs[4], resp.Data[i].ID, "deleted tenant should not appear") + for i := 0; i < 2 && i < len(resp.Models); i++ { + assert.NotEqual(t, tenantIDs[3], resp.Models[i].ID, "deleted tenant should not appear") + assert.NotEqual(t, tenantIDs[4], resp.Models[i].ID, "deleted tenant should not appear") } // Cleanup @@ -1370,10 +1376,10 @@ func (s *ListTenantTestSuite) TestListTenantKeysetPagination() { // Fetch page 1 (items 14, 13, 12, 11, 10 with DESC order) resp1, err := s.entityStore.ListTenant(s.ctx, models.ListTenantRequest{Limit: 5}) require.NoError(t, err) - require.Len(t, resp1.Data, 5, "page 1 should have 5 items") + require.Len(t, resp1.Models, 5, "page 1 should have 5 items") // Verify we got our test tenants - for _, tenant := range resp1.Data { + for _, tenant := range resp1.Models { require.Contains(t, tenant.ID, prefix, "page 1 should contain our test tenants") } @@ -1391,19 +1397,23 @@ func (s *ListTenantTestSuite) TestListTenantKeysetPagination() { // With keyset pagination, the cursor is the timestamp of item 10 // Page 2 will get items with timestamp < cursor, so items 9, 8, 7, 6, 5 // The new tenant has a newer timestamp so it won't appear on page 2 + var nextCursor string + if resp1.Pagination.Next != nil { + nextCursor = *resp1.Pagination.Next + } resp2, err := s.entityStore.ListTenant(s.ctx, models.ListTenantRequest{ Limit: 5, - Next: resp1.Next, + Next: nextCursor, }) require.NoError(t, err) - require.NotEmpty(t, resp2.Data, "page 2 should have items") + require.NotEmpty(t, resp2.Models, "page 2 should have items") // Verify no duplicates - first item on page 2 should NOT be the last from page 1 page1IDs := make(map[string]bool) - for _, tenant := range resp1.Data { + for _, tenant := range resp1.Models { page1IDs[tenant.ID] = true } - for _, tenant := range resp2.Data { + for _, tenant := range resp2.Models { assert.False(t, page1IDs[tenant.ID], "keyset pagination: no duplicates when adding during traversal, but found %s", tenant.ID) } @@ -1417,9 +1427,9 @@ func (s *ListTenantTestSuite) TestListTenantKeysetPagination() { // TestListTenantInputValidation tests input validation and error handling. func (s *ListTenantTestSuite) TestListTenantInputValidation() { - s.T().Run("invalid order returns error", func(t *testing.T) { + s.T().Run("invalid dir returns error", func(t *testing.T) { _, err := s.entityStore.ListTenant(s.ctx, models.ListTenantRequest{ - Order: "invalid", + Dir: "invalid", }) require.Error(t, err) assert.ErrorIs(t, err, models.ErrInvalidOrder) @@ -1466,7 +1476,7 @@ func (s *ListTenantTestSuite) TestListTenantInputValidation() { }) require.NoError(t, err) // Default limit is 20, should return up to 20 of the 25 tenants - assert.Equal(t, 20, len(resp.Data), "default limit should be 20") + assert.Equal(t, 20, len(resp.Models), "default limit should be 20") assert.Equal(t, 25, resp.Count, "total count should be 25") }) @@ -1488,22 +1498,22 @@ func (s *ListTenantTestSuite) TestListTenantInputValidation() { require.NoError(t, err) // Should succeed and return all 25 (capped to 100, but we only have 25) assert.NotNil(t, resp) - assert.Equal(t, 25, len(resp.Data), "should return all 25 tenants") + assert.Equal(t, 25, len(resp.Models), "should return all 25 tenants") assert.Equal(t, 25, resp.Count) }) - s.T().Run("empty order uses default desc", func(t *testing.T) { + s.T().Run("empty dir uses default desc", func(t *testing.T) { // Uses 25 tenants from SetupSuite created with sequential timestamps // s.tenants[0] is oldest, s.tenants[24] is newest resp, err := s.entityStore.ListTenant(s.ctx, models.ListTenantRequest{ - Order: "", // Should default to "desc" + Dir: "", // Should default to "desc" }) require.NoError(t, err) - require.Len(t, resp.Data, 20, "default limit is 20") + require.Len(t, resp.Models, 20, "default limit is 20") - // With DESC order, newest (s.tenants[24]) should be first + // With DESC dir, newest (s.tenants[24]) should be first // and older tenants should follow - assert.Equal(t, s.tenants[24].ID, resp.Data[0].ID, "newest tenant should be first with desc order") - assert.Equal(t, s.tenants[23].ID, resp.Data[1].ID, "second newest should be second") + assert.Equal(t, s.tenants[24].ID, resp.Models[0].ID, "newest tenant should be first with desc dir") + assert.Equal(t, s.tenants[23].ID, resp.Models[1].ID, "second newest should be second") }) } diff --git a/internal/models/tenant.go b/internal/models/tenant.go index a8f302c0e..0b30ef65b 100644 --- a/internal/models/tenant.go +++ b/internal/models/tenant.go @@ -15,25 +15,6 @@ type Tenant struct { UpdatedAt time.Time `json:"updated_at" redis:"updated_at"` } -// TenantListItem is a lightweight tenant representation for list operations. -// It excludes computed fields (destinations_count, topics) to avoid N+1 queries. -type TenantListItem struct { - ID string `json:"id"` - Metadata Metadata `json:"metadata,omitempty"` - CreatedAt time.Time `json:"created_at"` - UpdatedAt time.Time `json:"updated_at"` -} - -// ToListItem converts a Tenant to a TenantListItem. -func (t *Tenant) ToListItem() TenantListItem { - return TenantListItem{ - ID: t.ID, - Metadata: t.Metadata, - CreatedAt: t.CreatedAt, - UpdatedAt: t.UpdatedAt, - } -} - func (t *Tenant) parseRedisHash(hash map[string]string) error { if _, ok := hash["deleted_at"]; ok { return ErrTenantDeleted