Search: Build index from resource stats (#97320)
This commit is contained in:
@@ -121,6 +121,38 @@ func (b *backend) Stop(_ context.Context) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// GetResourceStats implements Backend.
|
||||
func (b *backend) GetResourceStats(ctx context.Context, minCount int) ([]resource.ResourceStats, error) {
|
||||
_, span := b.tracer.Start(ctx, tracePrefix+".GetResourceStats")
|
||||
defer span.End()
|
||||
|
||||
req := &sqlStatsRequest{
|
||||
SQLTemplate: sqltemplate.New(b.dialect),
|
||||
MinCount: minCount, // not used in query... yet?
|
||||
}
|
||||
|
||||
res := make([]resource.ResourceStats, 0, 100)
|
||||
err := b.db.WithTx(ctx, ReadCommittedRO, func(ctx context.Context, tx db.Tx) error {
|
||||
rows, err := dbutil.QueryRows(ctx, tx, sqlResourceStats, req)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
for rows.Next() {
|
||||
row := resource.ResourceStats{}
|
||||
err = rows.Scan(&row.Namespace, &row.Group, &row.Resource, &row.Count, &row.ResourceVersion)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if row.Count > int64(minCount) {
|
||||
res = append(res, row)
|
||||
}
|
||||
}
|
||||
return err
|
||||
})
|
||||
|
||||
return res, err
|
||||
}
|
||||
|
||||
func (b *backend) WriteEvent(ctx context.Context, event resource.WriteEvent) (int64, error) {
|
||||
_, span := b.tracer.Start(ctx, tracePrefix+"WriteEvent")
|
||||
defer span.End()
|
||||
@@ -137,35 +169,6 @@ func (b *backend) WriteEvent(ctx context.Context, event resource.WriteEvent) (in
|
||||
}
|
||||
}
|
||||
|
||||
// Namespaces returns the list of unique namespaces in storage.
|
||||
func (b *backend) Namespaces(ctx context.Context) ([]string, error) {
|
||||
var namespaces []string
|
||||
|
||||
err := b.db.WithTx(ctx, RepeatableRead, func(ctx context.Context, tx db.Tx) error {
|
||||
rows, err := tx.QueryContext(ctx, "SELECT DISTINCT(namespace) FROM resource ORDER BY namespace;")
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
defer func() {
|
||||
_ = rows.Close()
|
||||
}()
|
||||
|
||||
for rows.Next() {
|
||||
var ns string
|
||||
err = rows.Scan(&ns)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
namespaces = append(namespaces, ns)
|
||||
}
|
||||
|
||||
return nil
|
||||
})
|
||||
|
||||
return namespaces, err
|
||||
}
|
||||
|
||||
func (b *backend) create(ctx context.Context, event resource.WriteEvent) (int64, error) {
|
||||
ctx, span := b.tracer.Start(ctx, tracePrefix+"Create")
|
||||
defer span.End()
|
||||
|
||||
@@ -0,0 +1,12 @@
|
||||
SELECT
|
||||
{{ .Ident "namespace" }},
|
||||
{{ .Ident "group" }},
|
||||
{{ .Ident "resource" }},
|
||||
COUNT(*),
|
||||
MAX({{ .Ident "resource_version" }})
|
||||
FROM {{ .Ident "resource" }}
|
||||
GROUP BY
|
||||
{{ .Ident "namespace" }},
|
||||
{{ .Ident "group" }},
|
||||
{{ .Ident "resource" }}
|
||||
;
|
||||
@@ -31,6 +31,7 @@ var (
|
||||
sqlResourceInsert = mustTemplate("resource_insert.sql")
|
||||
sqlResourceUpdate = mustTemplate("resource_update.sql")
|
||||
sqlResourceRead = mustTemplate("resource_read.sql")
|
||||
sqlResourceStats = mustTemplate("resource_stats.sql")
|
||||
sqlResourceList = mustTemplate("resource_list.sql")
|
||||
sqlResourceHistoryList = mustTemplate("resource_history_list.sql")
|
||||
sqlResourceUpdateRV = mustTemplate("resource_update_rv.sql")
|
||||
@@ -71,6 +72,15 @@ func (r sqlResourceRequest) Validate() error {
|
||||
return nil // TODO
|
||||
}
|
||||
|
||||
type sqlStatsRequest struct {
|
||||
sqltemplate.SQLTemplate
|
||||
MinCount int
|
||||
}
|
||||
|
||||
func (r sqlStatsRequest) Validate() error {
|
||||
return nil // TODO
|
||||
}
|
||||
|
||||
type historyPollResponse struct {
|
||||
Key resource.ResourceKey
|
||||
ResourceVersion int64
|
||||
|
||||
@@ -219,5 +219,15 @@ func TestUnifiedStorageQueries(t *testing.T) {
|
||||
},
|
||||
},
|
||||
},
|
||||
|
||||
sqlResourceStats: {
|
||||
{
|
||||
Name: "query",
|
||||
Data: &sqlStatsRequest{
|
||||
SQLTemplate: mocks.NewTestingSQLTemplate(),
|
||||
MinCount: 10, // Not yet used in query (only response filter)
|
||||
},
|
||||
},
|
||||
},
|
||||
}})
|
||||
}
|
||||
|
||||
@@ -65,6 +65,7 @@ func NewResourceServer(ctx context.Context, db infraDB.DB, cfg *setting.Cfg,
|
||||
}, tracer, reg),
|
||||
Resources: docs,
|
||||
WorkerThreads: 5, // from cfg?
|
||||
InitMinCount: 1,
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -13,7 +13,6 @@ import (
|
||||
|
||||
"github.com/grafana/authlib/claims"
|
||||
"github.com/grafana/dskit/services"
|
||||
|
||||
"github.com/grafana/grafana/pkg/apimachinery/identity"
|
||||
"github.com/grafana/grafana/pkg/apimachinery/utils"
|
||||
infraDB "github.com/grafana/grafana/pkg/infra/db"
|
||||
@@ -98,6 +97,12 @@ func TestIntegrationBackendHappyPath(t *testing.T) {
|
||||
rv3, err = writeEvent(ctx, backend, "item3", resource.WatchEvent_ADDED)
|
||||
require.NoError(t, err)
|
||||
require.Greater(t, rv3, rv2)
|
||||
|
||||
stats, err := backend.GetResourceStats(ctx, 0)
|
||||
require.NoError(t, err)
|
||||
require.Len(t, stats, 1)
|
||||
require.Equal(t, int64(3), stats[0].Count)
|
||||
require.Equal(t, rv3, stats[0].ResourceVersion)
|
||||
})
|
||||
|
||||
t.Run("Update item2", func(t *testing.T) {
|
||||
|
||||
+12
@@ -0,0 +1,12 @@
|
||||
SELECT
|
||||
`namespace`,
|
||||
`group`,
|
||||
`resource`,
|
||||
COUNT(*),
|
||||
MAX(`resource_version`)
|
||||
FROM `resource`
|
||||
GROUP BY
|
||||
`namespace`,
|
||||
`group`,
|
||||
`resource`
|
||||
;
|
||||
+12
@@ -0,0 +1,12 @@
|
||||
SELECT
|
||||
"namespace",
|
||||
"group",
|
||||
"resource",
|
||||
COUNT(*),
|
||||
MAX("resource_version")
|
||||
FROM "resource"
|
||||
GROUP BY
|
||||
"namespace",
|
||||
"group",
|
||||
"resource"
|
||||
;
|
||||
+12
@@ -0,0 +1,12 @@
|
||||
SELECT
|
||||
"namespace",
|
||||
"group",
|
||||
"resource",
|
||||
COUNT(*),
|
||||
MAX("resource_version")
|
||||
FROM "resource"
|
||||
GROUP BY
|
||||
"namespace",
|
||||
"group",
|
||||
"resource"
|
||||
;
|
||||
Reference in New Issue
Block a user