Storage: Support continue at specified resource version (#84868)
* support continue at specified resource version * detect whether list continue pages need to use entity_history, remove BatchRead, expand selectQuery helper * refactor continue token handling * fix tests, increase history chunk size * lint fix
This commit is contained in:
@@ -21,57 +21,115 @@ func (d Direction) String() string {
|
||||
return "ASC"
|
||||
}
|
||||
|
||||
type joinQuery struct {
|
||||
query string
|
||||
args []any
|
||||
}
|
||||
|
||||
type whereClause struct {
|
||||
query string
|
||||
args []any
|
||||
}
|
||||
|
||||
type selectQuery struct {
|
||||
dialect migrator.Dialect
|
||||
fields []string // SELECT xyz
|
||||
from string // FROM object
|
||||
fields []string // SELECT xyz
|
||||
from string // FROM object
|
||||
joins []joinQuery // JOIN object
|
||||
offset int64
|
||||
limit int64
|
||||
oneExtra bool
|
||||
|
||||
where []string
|
||||
args []any
|
||||
where []whereClause
|
||||
|
||||
groupBy []string
|
||||
|
||||
orderBy []string
|
||||
direction []Direction
|
||||
}
|
||||
|
||||
func (q *selectQuery) addWhere(f string, val ...any) {
|
||||
q.args = append(q.args, val...)
|
||||
// if the field contains a question mark, we assume it's a raw where clause
|
||||
if strings.Contains(f, "?") {
|
||||
q.where = append(q.where, f)
|
||||
// otherwise we assume it's a field name
|
||||
} else {
|
||||
q.where = append(q.where, q.dialect.Quote(f)+"=?")
|
||||
func NewSelectQuery(dialect migrator.Dialect, from string) *selectQuery {
|
||||
return &selectQuery{
|
||||
dialect: dialect,
|
||||
from: from,
|
||||
}
|
||||
}
|
||||
|
||||
func (q *selectQuery) addWhereInSubquery(f string, subquery string, subqueryArgs []any) {
|
||||
q.args = append(q.args, subqueryArgs...)
|
||||
q.where = append(q.where, q.dialect.Quote(f)+" IN ("+subquery+")")
|
||||
func (q *selectQuery) From(from string) {
|
||||
q.from = from
|
||||
}
|
||||
|
||||
func (q *selectQuery) addWhereIn(f string, vals []string) {
|
||||
func (q *selectQuery) SetLimit(limit int64) {
|
||||
q.limit = limit
|
||||
}
|
||||
|
||||
func (q *selectQuery) SetOffset(offset int64) {
|
||||
q.offset = offset
|
||||
}
|
||||
|
||||
func (q *selectQuery) SetOneExtra() {
|
||||
q.oneExtra = true
|
||||
}
|
||||
|
||||
func (q *selectQuery) UnsetOneExtra() {
|
||||
q.oneExtra = false
|
||||
}
|
||||
|
||||
func (q *selectQuery) AddFields(f ...string) {
|
||||
for _, field := range f {
|
||||
q.fields = append(q.fields, "t."+q.dialect.Quote(field))
|
||||
}
|
||||
}
|
||||
|
||||
func (q *selectQuery) AddRawFields(f ...string) {
|
||||
q.fields = append(q.fields, f...)
|
||||
}
|
||||
|
||||
func (q *selectQuery) AddJoin(j string, args ...any) {
|
||||
q.joins = append(q.joins, joinQuery{query: j, args: args})
|
||||
}
|
||||
|
||||
func (q *selectQuery) AddWhere(f string, val ...any) {
|
||||
// if the field contains a question mark, we assume it's a raw where clause
|
||||
if strings.Contains(f, "?") {
|
||||
q.where = append(q.where, whereClause{f, val})
|
||||
// otherwise we assume it's a field name
|
||||
} else {
|
||||
q.where = append(q.where, whereClause{"t." + q.dialect.Quote(f) + "=?", val})
|
||||
}
|
||||
}
|
||||
|
||||
func (q *selectQuery) AddWhereInSubquery(f string, subquery string, subqueryArgs []any) {
|
||||
q.where = append(q.where, whereClause{"t." + q.dialect.Quote(f) + " IN (" + subquery + ")", subqueryArgs})
|
||||
}
|
||||
|
||||
func (q *selectQuery) AddWhereIn(f string, vals []any) {
|
||||
count := len(vals)
|
||||
if count > 1 {
|
||||
sb := strings.Builder{}
|
||||
sb.WriteString(q.dialect.Quote(f))
|
||||
sb.WriteString("t." + q.dialect.Quote(f))
|
||||
sb.WriteString(" IN (")
|
||||
for i := 0; i < count; i++ {
|
||||
if i > 0 {
|
||||
sb.WriteString(",")
|
||||
}
|
||||
sb.WriteString("?")
|
||||
q.args = append(q.args, vals[i])
|
||||
}
|
||||
sb.WriteString(") ")
|
||||
q.where = append(q.where, sb.String())
|
||||
q.where = append(q.where, whereClause{sb.String(), vals})
|
||||
} else if count == 1 {
|
||||
q.addWhere(f, vals[0])
|
||||
q.AddWhere(f, vals[0])
|
||||
}
|
||||
}
|
||||
|
||||
func ToAnyList[T any](input []T) []any {
|
||||
list := make([]any, len(input))
|
||||
for i, v := range input {
|
||||
list[i] = v
|
||||
}
|
||||
return list
|
||||
}
|
||||
|
||||
const sqlLikeEscape = "#"
|
||||
|
||||
var sqlLikeEscapeReplacer = strings.NewReplacer(
|
||||
@@ -85,39 +143,57 @@ func escapeJSONStringSQLLike(s string) string {
|
||||
return sqlLikeEscapeReplacer.Replace(string(b))
|
||||
}
|
||||
|
||||
func (q *selectQuery) addWhereJsonContainsKV(field string, key string, value string) {
|
||||
func (q *selectQuery) AddWhereJsonContainsKV(field string, key string, value string) {
|
||||
escapedKey := escapeJSONStringSQLLike(key)
|
||||
escapedValue := escapeJSONStringSQLLike(value)
|
||||
q.where = append(q.where, q.dialect.Quote(field)+" LIKE ? ESCAPE ?")
|
||||
q.args = append(q.args, "{%"+escapedKey+":"+escapedValue+"%}", sqlLikeEscape)
|
||||
q.where = append(q.where, whereClause{
|
||||
"t." + q.dialect.Quote(field) + " LIKE ? ESCAPE ?",
|
||||
[]any{"{%\"" + escapedKey + "\":\"" + escapedValue + "\"%}", sqlLikeEscape},
|
||||
})
|
||||
}
|
||||
|
||||
func (q *selectQuery) addOrderBy(field string, direction Direction) {
|
||||
func (q *selectQuery) AddGroupBy(f string) {
|
||||
q.groupBy = append(q.groupBy, f)
|
||||
}
|
||||
|
||||
func (q *selectQuery) AddOrderBy(field string, direction Direction) {
|
||||
q.orderBy = append(q.orderBy, field)
|
||||
q.direction = append(q.direction, direction)
|
||||
}
|
||||
|
||||
func (q *selectQuery) toQuery() (string, []any) {
|
||||
args := q.args
|
||||
func (q *selectQuery) ToQuery() (string, []any) {
|
||||
args := []any{}
|
||||
sb := strings.Builder{}
|
||||
sb.WriteString("SELECT ")
|
||||
quotedFields := make([]string, len(q.fields))
|
||||
for i, f := range q.fields {
|
||||
quotedFields[i] = q.dialect.Quote(f)
|
||||
}
|
||||
sb.WriteString(strings.Join(quotedFields, ","))
|
||||
sb.WriteString(strings.Join(q.fields, ","))
|
||||
sb.WriteString(" FROM ")
|
||||
sb.WriteString(q.from)
|
||||
sb.WriteString(" AS t")
|
||||
|
||||
for _, j := range q.joins {
|
||||
sb.WriteString(" " + j.query)
|
||||
args = append(args, j.args...)
|
||||
}
|
||||
|
||||
// Templated where string
|
||||
where := len(q.where)
|
||||
if where > 0 {
|
||||
if len(q.where) > 0 {
|
||||
sb.WriteString(" WHERE ")
|
||||
for i := 0; i < where; i++ {
|
||||
for i, w := range q.where {
|
||||
if i > 0 {
|
||||
sb.WriteString(" AND ")
|
||||
}
|
||||
sb.WriteString(q.where[i])
|
||||
sb.WriteString(w.query)
|
||||
args = append(args, w.args...)
|
||||
}
|
||||
}
|
||||
|
||||
if len(q.groupBy) > 0 {
|
||||
sb.WriteString(" GROUP BY ")
|
||||
for i, f := range q.groupBy {
|
||||
if i > 0 {
|
||||
sb.WriteString(",")
|
||||
}
|
||||
sb.WriteString("t." + q.dialect.Quote(f))
|
||||
}
|
||||
}
|
||||
|
||||
@@ -127,21 +203,19 @@ func (q *selectQuery) toQuery() (string, []any) {
|
||||
if i > 0 {
|
||||
sb.WriteString(",")
|
||||
}
|
||||
sb.WriteString(q.dialect.Quote(f))
|
||||
sb.WriteString("t." + q.dialect.Quote(f))
|
||||
sb.WriteString(" ")
|
||||
sb.WriteString(q.direction[i].String())
|
||||
}
|
||||
}
|
||||
|
||||
limit := q.limit
|
||||
if limit < 1 {
|
||||
limit = 20
|
||||
q.limit = limit
|
||||
if limit > 0 {
|
||||
if q.oneExtra {
|
||||
limit = limit + 1
|
||||
}
|
||||
sb.WriteString(q.dialect.LimitOffset(limit, q.offset))
|
||||
}
|
||||
if q.oneExtra {
|
||||
limit = limit + 1
|
||||
}
|
||||
sb.WriteString(q.dialect.LimitOffset(limit, q.offset))
|
||||
|
||||
return sb.String(), args
|
||||
}
|
||||
|
||||
@@ -25,6 +25,9 @@ import (
|
||||
"github.com/grafana/grafana/pkg/services/store/entity/db"
|
||||
)
|
||||
|
||||
const entityTable = "entity"
|
||||
const entityHistoryTable = "entity_history"
|
||||
|
||||
// Make sure we implement both store + admin
|
||||
var _ entity.EntityStoreServer = &sqlEntityServer{}
|
||||
|
||||
@@ -203,7 +206,7 @@ func (s *sqlEntityServer) Read(ctx context.Context, r *entity.ReadEntityRequest)
|
||||
}
|
||||
|
||||
func (s *sqlEntityServer) read(ctx context.Context, tx session.SessionQuerier, r *entity.ReadEntityRequest) (*entity.Entity, error) {
|
||||
table := "entity"
|
||||
table := entityTable
|
||||
where := []string{}
|
||||
args := []any{}
|
||||
|
||||
@@ -220,7 +223,7 @@ func (s *sqlEntityServer) read(ctx context.Context, tx session.SessionQuerier, r
|
||||
args = append(args, key.Namespace, key.Group, key.Resource, key.Name)
|
||||
|
||||
if r.ResourceVersion != 0 {
|
||||
table = "entity_history"
|
||||
table = entityHistoryTable
|
||||
where = append(where, s.dialect.Quote("resource_version")+">=?")
|
||||
args = append(args, r.ResourceVersion)
|
||||
}
|
||||
@@ -257,58 +260,6 @@ func (s *sqlEntityServer) read(ctx context.Context, tx session.SessionQuerier, r
|
||||
return rowToEntity(rows, r)
|
||||
}
|
||||
|
||||
func (s *sqlEntityServer) BatchRead(ctx context.Context, b *entity.BatchReadEntityRequest) (*entity.BatchReadEntityResponse, error) {
|
||||
if len(b.Batch) < 1 {
|
||||
return nil, fmt.Errorf("missing querires")
|
||||
}
|
||||
|
||||
first := b.Batch[0]
|
||||
args := []any{}
|
||||
constraints := []string{}
|
||||
|
||||
for _, r := range b.Batch {
|
||||
if r.WithBody != first.WithBody || r.WithStatus != first.WithStatus {
|
||||
return nil, fmt.Errorf("requests must want the same things")
|
||||
}
|
||||
|
||||
if r.Key == "" {
|
||||
return nil, fmt.Errorf("missing key")
|
||||
}
|
||||
|
||||
constraints = append(constraints, s.dialect.Quote("key")+"=?")
|
||||
args = append(args, r.Key)
|
||||
|
||||
if r.ResourceVersion != 0 {
|
||||
return nil, fmt.Errorf("version not supported for batch read (yet?)")
|
||||
}
|
||||
}
|
||||
|
||||
req := b.Batch[0]
|
||||
query, err := s.getReadSelect(req)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
query += " FROM entity" +
|
||||
" WHERE (" + strings.Join(constraints, " OR ") + ")"
|
||||
rows, err := s.sess.Query(ctx, query, args...)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer func() { _ = rows.Close() }()
|
||||
|
||||
// TODO? make sure the results are in order?
|
||||
rsp := &entity.BatchReadEntityResponse{}
|
||||
for rows.Next() {
|
||||
r, err := rowToEntity(rows, req)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
rsp.Results = append(rsp.Results, r)
|
||||
}
|
||||
return rsp, nil
|
||||
}
|
||||
|
||||
//nolint:gocyclo
|
||||
func (s *sqlEntityServer) Create(ctx context.Context, r *entity.CreateEntityRequest) (*entity.CreateEntityResponse, error) {
|
||||
if err := s.Init(); err != nil {
|
||||
@@ -488,13 +439,13 @@ func (s *sqlEntityServer) Create(ctx context.Context, r *entity.CreateEntityRequ
|
||||
}
|
||||
|
||||
// 1. Add row to the `entity_history` values
|
||||
if err := s.dialect.Insert(ctx, tx, "entity_history", values); err != nil {
|
||||
if err := s.dialect.Insert(ctx, tx, entityHistoryTable, values); err != nil {
|
||||
s.log.Error("error inserting entity history", "msg", err.Error())
|
||||
return err
|
||||
}
|
||||
|
||||
// 2. Add row to the main `entity` table
|
||||
if err := s.dialect.Insert(ctx, tx, "entity", values); err != nil {
|
||||
if err := s.dialect.Insert(ctx, tx, entityTable, values); err != nil {
|
||||
s.log.Error("error inserting entity", "msg", err.Error())
|
||||
return err
|
||||
}
|
||||
@@ -697,7 +648,7 @@ func (s *sqlEntityServer) Update(ctx context.Context, r *entity.UpdateEntityRequ
|
||||
}
|
||||
|
||||
// 1. Add the `entity_history` values
|
||||
if err := s.dialect.Insert(ctx, tx, "entity_history", values); err != nil {
|
||||
if err := s.dialect.Insert(ctx, tx, entityHistoryTable, values); err != nil {
|
||||
s.log.Error("error inserting entity history", "msg", err.Error())
|
||||
return err
|
||||
}
|
||||
@@ -717,7 +668,7 @@ func (s *sqlEntityServer) Update(ctx context.Context, r *entity.UpdateEntityRequ
|
||||
err = s.dialect.Update(
|
||||
ctx,
|
||||
tx,
|
||||
"entity",
|
||||
entityTable,
|
||||
values,
|
||||
map[string]any{
|
||||
"guid": current.Guid,
|
||||
@@ -898,7 +849,7 @@ func (s *sqlEntityServer) doDelete(ctx context.Context, tx *session.SessionTx, e
|
||||
}
|
||||
|
||||
// 1. Add the `entity_history` values
|
||||
if err := s.dialect.Insert(ctx, tx, "entity_history", values); err != nil {
|
||||
if err := s.dialect.Insert(ctx, tx, entityHistoryTable, values); err != nil {
|
||||
s.log.Error("error inserting entity history", "msg", err.Error())
|
||||
return err
|
||||
}
|
||||
@@ -962,28 +913,29 @@ func (s *sqlEntityServer) History(ctx context.Context, r *entity.EntityHistoryRe
|
||||
limit = r.Limit
|
||||
}
|
||||
|
||||
fields := s.getReadFields(r)
|
||||
|
||||
entityQuery := selectQuery{
|
||||
dialect: s.dialect,
|
||||
fields: fields,
|
||||
from: "entity_history", // the table
|
||||
args: []any{},
|
||||
from: entityHistoryTable, // the table
|
||||
limit: r.Limit,
|
||||
offset: 0,
|
||||
oneExtra: true, // request one more than the limit (and show next token if it exists)
|
||||
}
|
||||
|
||||
fields := s.getReadFields(r)
|
||||
entityQuery.AddFields(fields...)
|
||||
|
||||
args := []any{key.Group, key.Resource}
|
||||
whereclause := "(" + s.dialect.Quote("group") + "=? AND " + s.dialect.Quote("resource") + "=?"
|
||||
if key.Namespace != "" {
|
||||
args = append(args, key.Namespace)
|
||||
whereclause += " AND " + s.dialect.Quote("namespace") + "=?"
|
||||
}
|
||||
args = append(args, key.Name)
|
||||
whereclause += " AND " + s.dialect.Quote("name") + "=?)"
|
||||
if key.Name != "" {
|
||||
args = append(args, key.Name)
|
||||
whereclause += " AND " + s.dialect.Quote("name") + "=?"
|
||||
}
|
||||
whereclause += ")"
|
||||
|
||||
entityQuery.addWhere(whereclause, args...)
|
||||
entityQuery.AddWhere(whereclause, args...)
|
||||
|
||||
// if we have a page token, use that to specify the first record
|
||||
continueToken, err := GetContinueToken(r)
|
||||
@@ -999,11 +951,11 @@ func (s *sqlEntityServer) History(ctx context.Context, r *entity.EntityHistoryRe
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
entityQuery.addOrderBy(sortBy.Field, sortBy.Direction)
|
||||
entityQuery.AddOrderBy(sortBy.Field, sortBy.Direction)
|
||||
}
|
||||
entityQuery.addOrderBy("resource_version", Ascending)
|
||||
entityQuery.AddOrderBy("resource_version", Ascending)
|
||||
|
||||
query, args := entityQuery.toQuery()
|
||||
query, args := entityQuery.ToQuery()
|
||||
|
||||
s.log.Debug("history", "query", query, "args", args)
|
||||
|
||||
@@ -1044,8 +996,10 @@ type ContinueRequest interface {
|
||||
}
|
||||
|
||||
type ContinueToken struct {
|
||||
Sort []string `json:"s"`
|
||||
StartOffset int64 `json:"o"`
|
||||
Sort []string `json:"s"`
|
||||
StartOffset int64 `json:"o"`
|
||||
ResourceVersion int64 `json:"v"`
|
||||
RecordCnt int64 `json:"c"`
|
||||
}
|
||||
|
||||
func (c *ContinueToken) String() string {
|
||||
@@ -1112,6 +1066,7 @@ func ParseSortBy(sort string) (*SortBy, error) {
|
||||
return sortBy, nil
|
||||
}
|
||||
|
||||
//nolint:gocyclo
|
||||
func (s *sqlEntityServer) List(ctx context.Context, r *entity.EntityListRequest) (*entity.EntityListResponse, error) {
|
||||
if err := s.Init(); err != nil {
|
||||
return nil, err
|
||||
@@ -1127,31 +1082,44 @@ func (s *sqlEntityServer) List(ctx context.Context, r *entity.EntityListRequest)
|
||||
|
||||
fields := s.getReadFields(r)
|
||||
|
||||
entityQuery := selectQuery{
|
||||
dialect: s.dialect,
|
||||
fields: fields,
|
||||
from: "entity", // the table
|
||||
args: []any{},
|
||||
limit: r.Limit,
|
||||
offset: 0,
|
||||
oneExtra: true, // request one more than the limit (and show next token if it exists)
|
||||
}
|
||||
// main query we will use to retrieve entities
|
||||
entityQuery := NewSelectQuery(s.dialect, entityTable)
|
||||
entityQuery.AddFields(fields...)
|
||||
entityQuery.SetLimit(r.Limit)
|
||||
entityQuery.SetOneExtra()
|
||||
|
||||
// query to retrieve the max resource version and entity count
|
||||
rvMaxQuery := NewSelectQuery(s.dialect, entityTable)
|
||||
rvMaxQuery.AddRawFields("coalesce(max(resource_version),0) as rv", "count(guid) as cnt")
|
||||
|
||||
// subquery to get latest resource version for each entity
|
||||
// when we need to query from entity_history
|
||||
rvSubQuery := NewSelectQuery(s.dialect, entityHistoryTable)
|
||||
rvSubQuery.AddFields("guid")
|
||||
rvSubQuery.AddRawFields("max(resource_version) as max_rv")
|
||||
|
||||
// if we are looking for deleted entities, we list "deleted" entries from the entity_history table
|
||||
if r.Deleted {
|
||||
entityQuery.from = "entity_history"
|
||||
entityQuery.addWhere("action", entity.Entity_DELETED)
|
||||
entityQuery.from = entityHistoryTable
|
||||
entityQuery.AddWhere("action", entity.Entity_DELETED)
|
||||
|
||||
rvMaxQuery.from = entityHistoryTable
|
||||
rvMaxQuery.AddWhere("action", entity.Entity_DELETED)
|
||||
}
|
||||
|
||||
// TODO fix this
|
||||
// entityQuery.addWhere("namespace", user.OrgID)
|
||||
|
||||
if len(r.Group) > 0 {
|
||||
entityQuery.addWhereIn("group", r.Group)
|
||||
entityQuery.AddWhereIn("group", ToAnyList(r.Group))
|
||||
rvMaxQuery.AddWhereIn("group", ToAnyList(r.Group))
|
||||
rvSubQuery.AddWhereIn("group", ToAnyList(r.Group))
|
||||
}
|
||||
|
||||
if len(r.Resource) > 0 {
|
||||
entityQuery.addWhereIn("resource", r.Resource)
|
||||
entityQuery.AddWhereIn("resource", ToAnyList(r.Resource))
|
||||
rvMaxQuery.AddWhereIn("resource", ToAnyList(r.Resource))
|
||||
rvSubQuery.AddWhereIn("resource", ToAnyList(r.Resource))
|
||||
}
|
||||
|
||||
if len(r.Key) > 0 {
|
||||
@@ -1164,27 +1132,41 @@ func (s *sqlEntityServer) List(ctx context.Context, r *entity.EntityListRequest)
|
||||
}
|
||||
|
||||
args = append(args, key.Group, key.Resource)
|
||||
whereclause := "(" + s.dialect.Quote("group") + "=? AND " + s.dialect.Quote("resource") + "=?"
|
||||
whereclause := "(t." + s.dialect.Quote("group") + "=? AND t." + s.dialect.Quote("resource") + "=?"
|
||||
if key.Namespace != "" {
|
||||
args = append(args, key.Namespace)
|
||||
whereclause += " AND " + s.dialect.Quote("namespace") + "=?"
|
||||
whereclause += " AND t." + s.dialect.Quote("namespace") + "=?"
|
||||
}
|
||||
if key.Name != "" {
|
||||
args = append(args, key.Name)
|
||||
whereclause += " AND " + s.dialect.Quote("name") + "=?"
|
||||
whereclause += " AND t." + s.dialect.Quote("name") + "=?"
|
||||
}
|
||||
whereclause += ")"
|
||||
|
||||
where = append(where, whereclause)
|
||||
}
|
||||
|
||||
entityQuery.addWhere("("+strings.Join(where, " OR ")+")", args...)
|
||||
entityQuery.AddWhere("("+strings.Join(where, " OR ")+")", args...)
|
||||
rvMaxQuery.AddWhere("("+strings.Join(where, " OR ")+")", args...)
|
||||
rvSubQuery.AddWhere("("+strings.Join(where, " OR ")+")", args...)
|
||||
}
|
||||
|
||||
// Folder guid
|
||||
if r.Folder != "" {
|
||||
entityQuery.addWhere("folder", r.Folder)
|
||||
// get the maximum resource version and count of entities
|
||||
type RVMaxRow struct {
|
||||
Rv int64 `db:"rv"`
|
||||
Cnt int64 `db:"cnt"`
|
||||
}
|
||||
rvMaxRow := &RVMaxRow{}
|
||||
query, args := rvMaxQuery.ToQuery()
|
||||
|
||||
err = s.sess.Get(ctx, rvMaxRow, query, args...)
|
||||
if err != nil {
|
||||
if !errors.Is(err, sql.ErrNoRows) {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
|
||||
s.log.Debug("getting max rv", "maxRv", rvMaxRow.Rv, "cnt", rvMaxRow.Cnt, "query", query, "args", args)
|
||||
|
||||
// if we have a page token, use that to specify the first record
|
||||
continueToken, err := GetContinueToken(r)
|
||||
@@ -1193,13 +1175,53 @@ func (s *sqlEntityServer) List(ctx context.Context, r *entity.EntityListRequest)
|
||||
}
|
||||
if continueToken != nil {
|
||||
entityQuery.offset = continueToken.StartOffset
|
||||
if continueToken.ResourceVersion > 0 {
|
||||
if r.Deleted {
|
||||
// if we're continuing, we need to list only revisions that are older than the given resource version
|
||||
entityQuery.AddWhere("resource_version <= ?", continueToken.ResourceVersion)
|
||||
} else {
|
||||
// cap versions considered by the per resource max version subquery
|
||||
rvSubQuery.AddWhere("resource_version <= ?", continueToken.ResourceVersion)
|
||||
}
|
||||
}
|
||||
|
||||
if (continueToken.ResourceVersion > 0 && continueToken.ResourceVersion != rvMaxRow.Rv) || (continueToken.RecordCnt > 0 && continueToken.RecordCnt != rvMaxRow.Cnt) {
|
||||
entityQuery.From(entityHistoryTable)
|
||||
entityQuery.AddWhere("t.action != ?", entity.Entity_DELETED)
|
||||
|
||||
rvSubQuery.AddGroupBy("guid")
|
||||
query, args = rvSubQuery.ToQuery()
|
||||
entityQuery.AddJoin("INNER JOIN ("+query+") rv ON rv.guid = t.guid AND rv.max_rv = t.resource_version", args...)
|
||||
}
|
||||
} else {
|
||||
continueToken = &ContinueToken{
|
||||
Sort: r.Sort,
|
||||
StartOffset: 0,
|
||||
ResourceVersion: rvMaxRow.Rv,
|
||||
RecordCnt: rvMaxRow.Cnt,
|
||||
}
|
||||
|
||||
if continueToken.ResourceVersion == 0 {
|
||||
// we use a snowflake as a fallback resource version
|
||||
continueToken.ResourceVersion = s.snowflake.Generate().Int64()
|
||||
}
|
||||
}
|
||||
|
||||
// initialize the result
|
||||
rsp := &entity.EntityListResponse{
|
||||
ResourceVersion: continueToken.ResourceVersion,
|
||||
}
|
||||
|
||||
// Folder guid
|
||||
if r.Folder != "" {
|
||||
entityQuery.AddWhere("folder", r.Folder)
|
||||
}
|
||||
|
||||
if len(r.Labels) > 0 {
|
||||
// if we are looking for deleted entities, we need to use the labels column
|
||||
if r.Deleted {
|
||||
if entityQuery.from == entityHistoryTable {
|
||||
for labelKey, labelValue := range r.Labels {
|
||||
entityQuery.addWhereJsonContainsKV("labels", labelKey, labelValue)
|
||||
entityQuery.AddWhereJsonContainsKV("labels", labelKey, labelValue)
|
||||
}
|
||||
// for active entities, we can use the entity_labels table
|
||||
} else {
|
||||
@@ -1216,7 +1238,7 @@ func (s *sqlEntityServer) List(ctx context.Context, r *entity.EntityListRequest)
|
||||
" HAVING COUNT(label) = ?"
|
||||
args = append(args, len(r.Labels))
|
||||
|
||||
entityQuery.addWhereInSubquery("guid", query, args)
|
||||
entityQuery.AddWhereInSubquery("guid", query, args)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1225,11 +1247,11 @@ func (s *sqlEntityServer) List(ctx context.Context, r *entity.EntityListRequest)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
entityQuery.addOrderBy(sortBy.Field, sortBy.Direction)
|
||||
entityQuery.AddOrderBy(sortBy.Field, sortBy.Direction)
|
||||
}
|
||||
entityQuery.addOrderBy("guid", Ascending)
|
||||
entityQuery.AddOrderBy("guid", Ascending)
|
||||
|
||||
query, args := entityQuery.toQuery()
|
||||
query, args = entityQuery.ToQuery()
|
||||
|
||||
s.log.Debug("listing", "query", query, "args", args)
|
||||
|
||||
@@ -1238,9 +1260,6 @@ func (s *sqlEntityServer) List(ctx context.Context, r *entity.EntityListRequest)
|
||||
return nil, err
|
||||
}
|
||||
defer func() { _ = rows.Close() }()
|
||||
rsp := &entity.EntityListResponse{
|
||||
ResourceVersion: s.snowflake.Generate().Int64(),
|
||||
}
|
||||
for rows.Next() {
|
||||
result, err := rowToEntity(rows, r)
|
||||
if err != nil {
|
||||
@@ -1248,11 +1267,8 @@ func (s *sqlEntityServer) List(ctx context.Context, r *entity.EntityListRequest)
|
||||
}
|
||||
|
||||
// found more than requested
|
||||
if int64(len(rsp.Results)) >= entityQuery.limit {
|
||||
continueToken := &ContinueToken{
|
||||
Sort: r.Sort,
|
||||
StartOffset: entityQuery.offset + entityQuery.limit,
|
||||
}
|
||||
if entityQuery.limit > 0 && int64(len(rsp.Results)) >= entityQuery.limit {
|
||||
continueToken.StartOffset = entityQuery.offset + entityQuery.limit
|
||||
rsp.NextPageToken = continueToken.String()
|
||||
break
|
||||
}
|
||||
@@ -1298,18 +1314,18 @@ func (s *sqlEntityServer) watchInit(r *entity.EntityWatchRequest, w entity.Entit
|
||||
|
||||
entityQuery := selectQuery{
|
||||
dialect: s.dialect,
|
||||
fields: fields,
|
||||
from: "entity", // the table
|
||||
args: []any{},
|
||||
limit: 100, // r.Limit,
|
||||
oneExtra: true, // request one more than the limit (and show next token if it exists)
|
||||
from: entityTable, // the table
|
||||
limit: 1000, // r.Limit,
|
||||
oneExtra: true, // request one more than the limit (and show next token if it exists)
|
||||
}
|
||||
|
||||
entityQuery.AddFields(fields...)
|
||||
|
||||
// if we got an initial resource version, start from that location in the history
|
||||
fromZero := true
|
||||
if r.Since > 0 {
|
||||
entityQuery.from = "entity_history"
|
||||
entityQuery.addWhere("resource_version > ?", r.Since)
|
||||
entityQuery.from = entityHistoryTable
|
||||
entityQuery.AddWhere("resource_version > ?", r.Since)
|
||||
fromZero = false
|
||||
}
|
||||
|
||||
@@ -1317,7 +1333,7 @@ func (s *sqlEntityServer) watchInit(r *entity.EntityWatchRequest, w entity.Entit
|
||||
// entityQuery.addWhere("namespace", user.OrgID)
|
||||
|
||||
if len(r.Resource) > 0 {
|
||||
entityQuery.addWhereIn("resource", r.Resource)
|
||||
entityQuery.AddWhereIn("resource", ToAnyList(r.Resource))
|
||||
}
|
||||
|
||||
if len(r.Key) > 0 {
|
||||
@@ -1344,18 +1360,18 @@ func (s *sqlEntityServer) watchInit(r *entity.EntityWatchRequest, w entity.Entit
|
||||
where = append(where, whereclause)
|
||||
}
|
||||
|
||||
entityQuery.addWhere("("+strings.Join(where, " OR ")+")", args...)
|
||||
entityQuery.AddWhere("("+strings.Join(where, " OR ")+")", args...)
|
||||
}
|
||||
|
||||
// Folder guid
|
||||
if r.Folder != "" {
|
||||
entityQuery.addWhere("folder", r.Folder)
|
||||
entityQuery.AddWhere("folder", r.Folder)
|
||||
}
|
||||
|
||||
if len(r.Labels) > 0 {
|
||||
if r.Since > 0 {
|
||||
for labelKey, labelValue := range r.Labels {
|
||||
entityQuery.addWhereJsonContainsKV("labels", labelKey, labelValue)
|
||||
entityQuery.AddWhereJsonContainsKV("labels", labelKey, labelValue)
|
||||
}
|
||||
} else {
|
||||
var args []any
|
||||
@@ -1371,17 +1387,17 @@ func (s *sqlEntityServer) watchInit(r *entity.EntityWatchRequest, w entity.Entit
|
||||
" HAVING COUNT(label) = ?"
|
||||
args = append(args, len(r.Labels))
|
||||
|
||||
entityQuery.addWhereInSubquery("guid", query, args)
|
||||
entityQuery.AddWhereInSubquery("guid", query, args)
|
||||
}
|
||||
}
|
||||
|
||||
entityQuery.addOrderBy("resource_version", Ascending)
|
||||
entityQuery.AddOrderBy("resource_version", Ascending)
|
||||
|
||||
var err error
|
||||
|
||||
for hasmore := true; hasmore; {
|
||||
err = func() error {
|
||||
query, args := entityQuery.toQuery()
|
||||
query, args := entityQuery.ToQuery()
|
||||
|
||||
s.log.Debug("watch init", "query", query, "args", args)
|
||||
|
||||
@@ -1461,18 +1477,17 @@ func (s *sqlEntityServer) poll(since int64, out chan *entity.Entity) (int64, err
|
||||
err := func() error {
|
||||
entityQuery := selectQuery{
|
||||
dialect: s.dialect,
|
||||
fields: fields,
|
||||
from: "entity_history", // the table
|
||||
args: []any{},
|
||||
limit: 100, // r.Limit,
|
||||
from: entityHistoryTable, // the table
|
||||
limit: 100, // r.Limit,
|
||||
// offset: 0,
|
||||
oneExtra: true, // request one more than the limit (and show next token if it exists)
|
||||
orderBy: []string{"resource_version"},
|
||||
}
|
||||
|
||||
entityQuery.addWhere("resource_version > ?", since)
|
||||
entityQuery.AddFields(fields...)
|
||||
entityQuery.AddWhere("resource_version > ?", since)
|
||||
|
||||
query, args := entityQuery.toQuery()
|
||||
query, args := entityQuery.ToQuery()
|
||||
|
||||
rows, err := s.sess.Query(s.ctx, query, args...)
|
||||
if err != nil {
|
||||
|
||||
Reference in New Issue
Block a user