openmeter / streaming /clickhouse /event_query_v2_test.go
Leon4gr45's picture
Upload folder using huggingface_hub (part 6)
d6f631f verified
Raw
History Blame Contribute Delete
11.4 kB
package clickhouse
import (
"testing"
"time"
"github.com/samber/lo"
"github.com/stretchr/testify/assert"
"github.com/openmeterio/openmeter/openmeter/customer"
"github.com/openmeterio/openmeter/openmeter/streaming"
"github.com/openmeterio/openmeter/pkg/filter"
"github.com/openmeterio/openmeter/pkg/models"
"github.com/openmeterio/openmeter/pkg/pagination/v2"
"github.com/openmeterio/openmeter/pkg/sortx"
)
func TestQueryEventsTableV2_ToSQL(t *testing.T) {
now := time.Now()
limit := 50
cursorTime := time.Date(2024, 1, 1, 0, 0, 0, 0, time.UTC)
cursorID := "event-123"
tests := []struct {
name string
query queryEventsTableV2
wantSQL string
wantArgs []interface{}
}{
{
name: "basic query with namespace only",
query: queryEventsTableV2{
Database: "openmeter",
EventsTableName: "om_events",
Params: streaming.ListEventsV2Params{
Namespace: "my_namespace",
},
},
wantSQL: "SELECT id, type, subject, source, time, data, ingested_at, stored_at, store_row_id FROM openmeter.om_events WHERE namespace = ? ORDER BY time DESC, store_row_id DESC LIMIT ?",
wantArgs: []interface{}{"my_namespace", 100},
},
{
name: "query with ID filter",
query: queryEventsTableV2{
Database: "openmeter",
EventsTableName: "om_events",
Params: streaming.ListEventsV2Params{
Namespace: "my_namespace",
ID: &filter.FilterString{
Eq: lo.ToPtr("event-123"),
},
},
},
wantSQL: "SELECT id, type, subject, source, time, data, ingested_at, stored_at, store_row_id FROM openmeter.om_events WHERE namespace = ? AND id = ? ORDER BY time DESC, store_row_id DESC LIMIT ?",
wantArgs: []interface{}{"my_namespace", "event-123", 100},
},
{
name: "query with subject filter",
query: queryEventsTableV2{
Database: "openmeter",
EventsTableName: "om_events",
Params: streaming.ListEventsV2Params{
Namespace: "my_namespace",
Subject: &filter.FilterString{
Like: lo.ToPtr("%customer%"),
},
},
},
wantSQL: "SELECT id, type, subject, source, time, data, ingested_at, stored_at, store_row_id FROM openmeter.om_events WHERE namespace = ? AND subject LIKE ? ORDER BY time DESC, store_row_id DESC LIMIT ?",
wantArgs: []interface{}{"my_namespace", "%customer%", 100},
},
{
name: "query with time filter",
query: queryEventsTableV2{
Database: "openmeter",
EventsTableName: "om_events",
Params: streaming.ListEventsV2Params{
Namespace: "my_namespace",
Time: &filter.FilterTime{
Gte: &now,
},
},
},
wantSQL: "SELECT id, type, subject, source, time, data, ingested_at, stored_at, store_row_id FROM openmeter.om_events WHERE namespace = ? AND time >= ? ORDER BY time DESC, store_row_id DESC LIMIT ?",
wantArgs: []interface{}{"my_namespace", now, 100},
},
{
name: "query with cursor and custom limit",
query: queryEventsTableV2{
Database: "openmeter",
EventsTableName: "om_events",
Params: streaming.ListEventsV2Params{
Namespace: "my_namespace",
Cursor: &pagination.Cursor{
Time: cursorTime,
ID: cursorID,
},
Limit: &limit,
},
},
wantSQL: "SELECT id, type, subject, source, time, data, ingested_at, stored_at, store_row_id FROM openmeter.om_events WHERE namespace = ? AND time <= ? AND (time < ? OR store_row_id < ?) ORDER BY time DESC, store_row_id DESC LIMIT ?",
wantArgs: []interface{}{"my_namespace", cursorTime.Unix(), cursorTime.Unix(), cursorID, 50},
},
{
name: "query with ingested_at filter defaults to time ordering",
query: queryEventsTableV2{
Database: "openmeter",
EventsTableName: "om_events",
Params: streaming.ListEventsV2Params{
Namespace: "my_namespace",
IngestedAt: &filter.FilterTime{
Gte: &now,
},
},
},
wantSQL: "SELECT id, type, subject, source, time, data, ingested_at, stored_at, store_row_id FROM openmeter.om_events WHERE namespace = ? AND ingested_at >= ? ORDER BY time DESC, store_row_id DESC LIMIT ?",
wantArgs: []interface{}{"my_namespace", now, 100},
},
{
name: "query with explicit ingested_at sort",
query: queryEventsTableV2{
Database: "openmeter",
EventsTableName: "om_events",
Params: streaming.ListEventsV2Params{
Namespace: "my_namespace",
SortBy: streaming.EventSortFieldIngestedAt,
},
},
wantSQL: "SELECT id, type, subject, source, time, data, ingested_at, stored_at, store_row_id FROM openmeter.om_events WHERE namespace = ? ORDER BY ingested_at DESC, store_row_id DESC LIMIT ?",
wantArgs: []interface{}{"my_namespace", 100},
},
{
name: "query with stored_at filter",
query: queryEventsTableV2{
Database: "openmeter",
EventsTableName: "om_events",
Params: streaming.ListEventsV2Params{
Namespace: "my_namespace",
StoredAt: &filter.FilterTime{
Gte: &now,
},
},
},
wantSQL: "SELECT id, type, subject, source, time, data, ingested_at, stored_at, store_row_id FROM openmeter.om_events WHERE namespace = ? AND stored_at >= ? ORDER BY time DESC, store_row_id DESC LIMIT ?",
wantArgs: []interface{}{"my_namespace", now, 100},
},
{
name: "query with explicit stored_at sort",
query: queryEventsTableV2{
Database: "openmeter",
EventsTableName: "om_events",
Params: streaming.ListEventsV2Params{
Namespace: "my_namespace",
SortBy: streaming.EventSortFieldStoredAt,
},
},
wantSQL: "SELECT id, type, subject, source, time, data, ingested_at, stored_at, store_row_id FROM openmeter.om_events WHERE namespace = ? ORDER BY stored_at DESC, store_row_id DESC LIMIT ?",
wantArgs: []interface{}{"my_namespace", 100},
},
{
name: "query with ascending time sort and cursor",
query: queryEventsTableV2{
Database: "openmeter",
EventsTableName: "om_events",
Params: streaming.ListEventsV2Params{
Namespace: "my_namespace",
Cursor: &pagination.Cursor{
Time: cursorTime,
ID: cursorID,
},
SortOrder: sortx.OrderAsc,
},
},
wantSQL: "SELECT id, type, subject, source, time, data, ingested_at, stored_at, store_row_id FROM openmeter.om_events WHERE namespace = ? AND time >= ? AND (time > ? OR store_row_id > ?) ORDER BY time ASC, store_row_id ASC LIMIT ?",
wantArgs: []interface{}{"my_namespace", cursorTime.Unix(), cursorTime.Unix(), cursorID, 100},
},
{
name: "query with customer filter",
query: queryEventsTableV2{
Database: "openmeter",
EventsTableName: "om_events",
Params: streaming.ListEventsV2Params{
Namespace: "my_namespace",
Customers: &[]streaming.Customer{
customer.Customer{
ManagedResource: models.ManagedResource{
NamespacedModel: models.NamespacedModel{
Namespace: "my_namespace",
},
ID: "customer1-id",
},
Key: lo.ToPtr("customer1-key"),
UsageAttribution: &customer.CustomerUsageAttribution{
SubjectKeys: []string{"customer1-subject1", "customer1-subject2"},
},
},
customer.Customer{
ManagedResource: models.ManagedResource{
NamespacedModel: models.NamespacedModel{
Namespace: "my_namespace",
},
ID: "customer2-id",
},
Key: lo.ToPtr("customer2-key"),
UsageAttribution: &customer.CustomerUsageAttribution{
SubjectKeys: []string{"customer2-subject1", "customer2-subject2"},
},
},
},
},
},
wantSQL: "WITH map('customer1-key', 'customer1-id', 'customer1-subject1', 'customer1-id', 'customer1-subject2', 'customer1-id', 'customer2-key', 'customer2-id', 'customer2-subject1', 'customer2-id', 'customer2-subject2', 'customer2-id') as subject_to_customer_id SELECT id, type, subject, source, time, data, ingested_at, stored_at, store_row_id, subject_to_customer_id[om_events.subject] AS customer_id FROM openmeter.om_events WHERE namespace = ? AND openmeter.om_events.subject IN (?) ORDER BY time DESC, store_row_id DESC LIMIT ?",
wantArgs: []interface{}{"my_namespace", []string{"customer1-key", "customer1-subject1", "customer1-subject2", "customer2-key", "customer2-subject1", "customer2-subject2"}, 100},
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
gotSQL, gotArgs := tt.query.toSQL()
assert.Equal(t, tt.wantSQL, gotSQL)
assert.Equal(t, tt.wantArgs, gotArgs)
})
}
}
func TestQueryEventsTableV2_ToCountRowSQL(t *testing.T) {
now := time.Now()
tests := []struct {
name string
query queryEventsTableV2
wantSQL string
wantArgs []interface{}
}{
{
name: "basic count query",
query: queryEventsTableV2{
Database: "openmeter",
EventsTableName: "om_events",
Params: streaming.ListEventsV2Params{
Namespace: "my_namespace",
},
},
wantSQL: "SELECT count() as total FROM openmeter.om_events WHERE namespace = ?",
wantArgs: []interface{}{"my_namespace"},
},
{
name: "count query with type filter",
query: queryEventsTableV2{
Database: "openmeter",
EventsTableName: "om_events",
Params: streaming.ListEventsV2Params{
Namespace: "my_namespace",
Type: &filter.FilterString{
Eq: lo.ToPtr("api-calls"),
},
},
},
wantSQL: "SELECT count() as total FROM openmeter.om_events WHERE namespace = ? AND type = ?",
wantArgs: []interface{}{"my_namespace", "api-calls"},
},
{
name: "count query with time filter",
query: queryEventsTableV2{
Database: "openmeter",
EventsTableName: "om_events",
Params: streaming.ListEventsV2Params{
Namespace: "my_namespace",
Time: &filter.FilterTime{
Gte: &now,
},
},
},
wantSQL: "SELECT count() as total FROM openmeter.om_events WHERE namespace = ? AND time >= ?",
wantArgs: []interface{}{"my_namespace", now},
},
{
name: "count query with subject filter",
query: queryEventsTableV2{
Database: "openmeter",
EventsTableName: "om_events",
Params: streaming.ListEventsV2Params{
Namespace: "my_namespace",
Subject: &filter.FilterString{
Like: lo.ToPtr("%customer%"),
},
},
},
wantSQL: "SELECT count() as total FROM openmeter.om_events WHERE namespace = ? AND subject LIKE ?",
wantArgs: []interface{}{"my_namespace", "%customer%"},
},
{
name: "count query with customer filter",
query: queryEventsTableV2{
Database: "openmeter",
EventsTableName: "om_events",
Params: streaming.ListEventsV2Params{
Namespace: "my_namespace",
Customers: &[]streaming.Customer{
customer.Customer{
ManagedResource: models.ManagedResource{
NamespacedModel: models.NamespacedModel{
Namespace: "my_namespace",
},
ID: "customer1",
},
UsageAttribution: &customer.CustomerUsageAttribution{
SubjectKeys: []string{"subject1"},
},
},
},
},
},
wantSQL: "SELECT count() as total FROM openmeter.om_events WHERE namespace = ? AND openmeter.om_events.subject IN (?)",
wantArgs: []interface{}{"my_namespace", []string{"subject1"}},
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
gotSQL, gotArgs := tt.query.toCountRowSQL()
assert.Equal(t, tt.wantSQL, gotSQL)
assert.Equal(t, tt.wantArgs, gotArgs)
})
}
}