| package clickhouse |
|
|
| import ( |
| "github.com/huandu/go-sqlbuilder" |
| "github.com/samber/lo" |
|
|
| "github.com/openmeterio/openmeter/openmeter/streaming" |
| "github.com/openmeterio/openmeter/pkg/sortx" |
| ) |
|
|
| const eventQueryV2DefaultLimit = 100 |
|
|
| |
| type queryEventsTableV2 struct { |
| Database string |
| EventsTableName string |
| Params streaming.ListEventsV2Params |
| } |
|
|
| |
| func (q queryEventsTableV2) toSQL() (string, []interface{}) { |
| tableName := getTableName(q.Database, q.EventsTableName) |
|
|
| query := sqlbuilder.ClickHouse.NewSelectBuilder() |
| query.Select("id", "type", "subject", "source", "time", "data", "ingested_at", "stored_at", "store_row_id") |
|
|
| |
| if q.Params.Customers != nil { |
| query = selectCustomerIdColumn(q.EventsTableName, *q.Params.Customers, query) |
| } |
|
|
| query.From(tableName) |
|
|
| |
| query.Where(query.Equal("namespace", q.Params.Namespace)) |
|
|
| if q.Params.ID != nil { |
| expr := q.Params.ID.SelectWhereExpr("id", query) |
| if expr != "" { |
| query.Where(expr) |
| } |
| } |
|
|
| if q.Params.Source != nil { |
| expr := q.Params.Source.SelectWhereExpr("source", query) |
| if expr != "" { |
| query.Where(expr) |
| } |
| } |
|
|
| if q.Params.Subject != nil { |
| expr := q.Params.Subject.SelectWhereExpr("subject", query) |
| if expr != "" { |
| query.Where(expr) |
| } |
| } |
|
|
| if q.Params.Customers != nil { |
| query = customersWhere(tableName, *q.Params.Customers, query) |
| } |
|
|
| if q.Params.Type != nil { |
| expr := q.Params.Type.SelectWhereExpr("type", query) |
| if expr != "" { |
| query.Where(expr) |
| } |
| } |
|
|
| if q.Params.Time != nil { |
| expr := q.Params.Time.SelectWhereExpr("time", query) |
| if expr != "" { |
| query.Where(expr) |
| } |
| } |
|
|
| if q.Params.IngestedAt != nil { |
| expr := q.Params.IngestedAt.SelectWhereExpr("ingested_at", query) |
| if expr != "" { |
| query.Where(expr) |
| } |
| } |
|
|
| if q.Params.StoredAt != nil { |
| expr := q.Params.StoredAt.SelectWhereExpr("stored_at", query) |
| if expr != "" { |
| query.Where(expr) |
| } |
| } |
|
|
| sortCol := string(q.Params.SortBy) |
| if sortCol == "" { |
| sortCol = string(streaming.EventSortFieldTime) |
| } |
|
|
| if q.Params.Cursor != nil { |
| if q.Params.SortOrder == sortx.OrderAsc { |
| query.Where( |
| |
| query.GreaterEqualThan(sortCol, q.Params.Cursor.Time.Unix()), |
| |
| |
| query.Or( |
| query.GreaterThan(sortCol, q.Params.Cursor.Time.Unix()), |
| query.GreaterThan("store_row_id", q.Params.Cursor.ID), |
| ), |
| ) |
| } else { |
| query.Where( |
| |
| query.LessEqualThan(sortCol, q.Params.Cursor.Time.Unix()), |
| |
| |
| query.Or( |
| query.LessThan(sortCol, q.Params.Cursor.Time.Unix()), |
| query.LessThan("store_row_id", q.Params.Cursor.ID), |
| ), |
| ) |
| } |
| } |
|
|
| switch q.Params.SortOrder { |
| case sortx.OrderAsc: |
| query.OrderByAsc(sortCol).OrderByAsc("store_row_id") |
| case sortx.OrderDesc: |
| fallthrough |
| default: |
| query.OrderByDesc(sortCol).OrderByDesc("store_row_id") |
| } |
|
|
| |
| query.Limit(lo.FromPtrOr(q.Params.Limit, eventQueryV2DefaultLimit)) |
|
|
| return query.Build() |
| } |
|
|
| |
| func (q queryEventsTableV2) toCountRowSQL() (string, []interface{}) { |
| tableName := getTableName(q.Database, q.EventsTableName) |
|
|
| query := sqlbuilder.ClickHouse.NewSelectBuilder() |
| query.Select("count() as total") |
| query.From(tableName) |
|
|
| |
| query.Where(query.Equal("namespace", q.Params.Namespace)) |
|
|
| |
| |
|
|
| if q.Params.Type != nil { |
| expr := q.Params.Type.SelectWhereExpr("type", query) |
| if expr != "" { |
| query.Where(expr) |
| } |
| } |
|
|
| if q.Params.Subject != nil { |
| expr := q.Params.Subject.SelectWhereExpr("subject", query) |
| if expr != "" { |
| query.Where(expr) |
| } |
| } |
|
|
| if q.Params.Customers != nil { |
| query = customersWhere(tableName, *q.Params.Customers, query) |
| } |
|
|
| if q.Params.Time != nil { |
| expr := q.Params.Time.SelectWhereExpr("time", query) |
| if expr != "" { |
| query.Where(expr) |
| } |
| } |
|
|
| return query.Build() |
| } |
|
|