kvstore: merge the metadata store into the datastore (#110334)
* migrate eventstore to datastore * Add folder to event key * lint * lint * lint * lint * remove foundkye * refactor the Keys methods to move the Sort outside the ListKey method * remove bad import * fix missing params * lint * fix test * perf improvement
This commit is contained in:
@@ -11,7 +11,6 @@ import (
|
||||
"math/rand/v2"
|
||||
"net/http"
|
||||
"sort"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/bwmarrin/snowflake"
|
||||
@@ -36,7 +35,6 @@ type kvStorageBackend struct {
|
||||
snowflake *snowflake.Node
|
||||
kv KV
|
||||
dataStore *dataStore
|
||||
metaStore *metadataStore
|
||||
eventStore *eventStore
|
||||
notifier *notifier
|
||||
builder DocumentBuilder
|
||||
@@ -83,7 +81,6 @@ func NewKvStorageBackend(opts KvBackendOptions) (StorageBackend, error) {
|
||||
backend := &kvStorageBackend{
|
||||
kv: kv,
|
||||
dataStore: newDataStore(kv),
|
||||
metaStore: newMetadataStore(kv),
|
||||
eventStore: eventStore,
|
||||
notifier: newNotifier(eventStore, notifierOptions{}),
|
||||
snowflake: s,
|
||||
@@ -139,16 +136,14 @@ func (k *kvStorageBackend) pruneEvents(ctx context.Context, key PruningKey) erro
|
||||
return fmt.Errorf("invalid pruning key, all fields must be set: %+v", key)
|
||||
}
|
||||
|
||||
listKey := ListRequestKey{
|
||||
counter := 0
|
||||
// iterate over all keys for the resource and delete versions beyond the latest 20
|
||||
for datakey, err := range k.dataStore.Keys(ctx, ListRequestKey{
|
||||
Namespace: key.Namespace,
|
||||
Group: key.Group,
|
||||
Resource: key.Resource,
|
||||
Name: key.Name,
|
||||
Sort: SortOrderDesc,
|
||||
}
|
||||
counter := 0
|
||||
// iterate over all keys for the resource and delete versions beyond the latest 20
|
||||
for datakey, err := range k.dataStore.Keys(ctx, listKey) {
|
||||
}, SortOrderDesc) {
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -217,10 +212,10 @@ func (k *kvStorageBackend) WriteEvent(ctx context.Context, event WriteEvent) (in
|
||||
case resourcepb.WatchEvent_ADDED:
|
||||
action = DataActionCreated
|
||||
// Check if resource already exists for create operations
|
||||
_, err := k.metaStore.GetLatestResourceKey(ctx, MetaGetRequestKey{
|
||||
Namespace: event.Key.Namespace,
|
||||
_, err := k.dataStore.GetLatestResourceKey(ctx, GetRequestKey{
|
||||
Group: event.Key.Group,
|
||||
Resource: event.Key.Resource,
|
||||
Namespace: event.Key.Namespace,
|
||||
Name: event.Key.Name,
|
||||
})
|
||||
if err == nil {
|
||||
@@ -244,44 +239,20 @@ func (k *kvStorageBackend) WriteEvent(ctx context.Context, event WriteEvent) (in
|
||||
return 0, fmt.Errorf("object is nil")
|
||||
}
|
||||
|
||||
// Build the search document
|
||||
doc, err := k.builder.BuildDocument(ctx, event.Key, rv, event.Value)
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("failed to build document: %w", err)
|
||||
}
|
||||
|
||||
// Write the data
|
||||
err = k.dataStore.Save(ctx, DataKey{
|
||||
Namespace: event.Key.Namespace,
|
||||
err := k.dataStore.Save(ctx, DataKey{
|
||||
Group: event.Key.Group,
|
||||
Resource: event.Key.Resource,
|
||||
Namespace: event.Key.Namespace,
|
||||
Name: event.Key.Name,
|
||||
ResourceVersion: rv,
|
||||
Action: action,
|
||||
Folder: obj.GetFolder(),
|
||||
}, bytes.NewReader(event.Value))
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("failed to write data: %w", err)
|
||||
}
|
||||
|
||||
// Write metadata
|
||||
err = k.metaStore.Save(ctx, MetaDataObj{
|
||||
Key: MetaDataKey{
|
||||
Namespace: event.Key.Namespace,
|
||||
Group: event.Key.Group,
|
||||
Resource: event.Key.Resource,
|
||||
Name: event.Key.Name,
|
||||
ResourceVersion: rv,
|
||||
Action: action,
|
||||
Folder: obj.GetFolder(),
|
||||
},
|
||||
Value: MetaData{
|
||||
IndexableDocument: *doc,
|
||||
},
|
||||
})
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("failed to write metadata: %w", err)
|
||||
}
|
||||
|
||||
// Write event
|
||||
err = k.eventStore.Save(ctx, Event{
|
||||
Namespace: event.Key.Namespace,
|
||||
@@ -311,10 +282,10 @@ func (k *kvStorageBackend) ReadResource(ctx context.Context, req *resourcepb.Rea
|
||||
if req.Key == nil {
|
||||
return &BackendReadResponse{Error: &resourcepb.ErrorResult{Code: http.StatusBadRequest, Message: "missing key"}}
|
||||
}
|
||||
meta, err := k.metaStore.GetResourceKeyAtRevision(ctx, MetaGetRequestKey{
|
||||
Namespace: req.Key.Namespace,
|
||||
meta, err := k.dataStore.GetResourceKeyAtRevision(ctx, GetRequestKey{
|
||||
Group: req.Key.Group,
|
||||
Resource: req.Key.Resource,
|
||||
Namespace: req.Key.Namespace,
|
||||
Name: req.Key.Name,
|
||||
}, req.ResourceVersion)
|
||||
if errors.Is(err, ErrNotFound) {
|
||||
@@ -323,12 +294,13 @@ func (k *kvStorageBackend) ReadResource(ctx context.Context, req *resourcepb.Rea
|
||||
return &BackendReadResponse{Error: &resourcepb.ErrorResult{Code: http.StatusInternalServerError, Message: err.Error()}}
|
||||
}
|
||||
data, err := k.dataStore.Get(ctx, DataKey{
|
||||
Namespace: req.Key.Namespace,
|
||||
Group: req.Key.Group,
|
||||
Resource: req.Key.Resource,
|
||||
Namespace: req.Key.Namespace,
|
||||
Name: req.Key.Name,
|
||||
ResourceVersion: meta.ResourceVersion,
|
||||
Action: meta.Action,
|
||||
Folder: meta.Folder,
|
||||
})
|
||||
if err != nil || data == nil {
|
||||
return &BackendReadResponse{Error: &resourcepb.ErrorResult{Code: http.StatusInternalServerError, Message: err.Error()}}
|
||||
@@ -369,12 +341,12 @@ func (k *kvStorageBackend) ListIterator(ctx context.Context, req *resourcepb.Lis
|
||||
}
|
||||
|
||||
// Fetch the latest objects
|
||||
keys := make([]MetaDataKey, 0, min(defaultListBufferSize, req.Limit+1))
|
||||
keys := make([]DataKey, 0, min(defaultListBufferSize, req.Limit+1))
|
||||
idx := 0
|
||||
for metaKey, err := range k.metaStore.ListResourceKeysAtRevision(ctx, MetaListRequestKey{
|
||||
Namespace: req.Options.Key.Namespace,
|
||||
for dataKey, err := range k.dataStore.ListResourceKeysAtRevision(ctx, ListRequestKey{
|
||||
Group: req.Options.Key.Group,
|
||||
Resource: req.Options.Key.Resource,
|
||||
Namespace: req.Options.Key.Namespace,
|
||||
Name: req.Options.Key.Name,
|
||||
}, resourceVersion) {
|
||||
if err != nil {
|
||||
@@ -385,7 +357,7 @@ func (k *kvStorageBackend) ListIterator(ctx context.Context, req *resourcepb.Lis
|
||||
idx++
|
||||
continue
|
||||
}
|
||||
keys = append(keys, metaKey)
|
||||
keys = append(keys, dataKey)
|
||||
// Only fetch the first limit items + 1 to get the next token.
|
||||
if len(keys) >= int(req.Limit+1) {
|
||||
break
|
||||
@@ -411,7 +383,7 @@ func (k *kvStorageBackend) ListIterator(ctx context.Context, req *resourcepb.Lis
|
||||
// kvListIterator implements ListIterator for KV storage
|
||||
type kvListIterator struct {
|
||||
ctx context.Context
|
||||
keys []MetaDataKey
|
||||
keys []DataKey
|
||||
currentIndex int
|
||||
dataStore *dataStore
|
||||
listRV int64
|
||||
@@ -437,14 +409,7 @@ func (i *kvListIterator) Next() bool {
|
||||
|
||||
i.rv, i.err = i.keys[i.currentIndex].ResourceVersion, nil
|
||||
|
||||
data, err := i.dataStore.Get(i.ctx, DataKey{
|
||||
Namespace: i.keys[i.currentIndex].Namespace,
|
||||
Group: i.keys[i.currentIndex].Group,
|
||||
Resource: i.keys[i.currentIndex].Resource,
|
||||
Name: i.keys[i.currentIndex].Name,
|
||||
ResourceVersion: i.keys[i.currentIndex].ResourceVersion,
|
||||
Action: i.keys[i.currentIndex].Action,
|
||||
})
|
||||
data, err := i.dataStore.Get(i.ctx, i.keys[i.currentIndex])
|
||||
if err != nil {
|
||||
i.err = err
|
||||
return false
|
||||
@@ -519,7 +484,7 @@ func filterHistoryKeysByVersion(historyKeys []DataKey, req *resourcepb.ListReque
|
||||
if req.ResourceVersion <= 0 {
|
||||
return nil, fmt.Errorf("expecting an explicit resource version query when using Exact matching")
|
||||
}
|
||||
var exactKeys []DataKey
|
||||
exactKeys := make([]DataKey, 0, len(historyKeys))
|
||||
for _, key := range historyKeys {
|
||||
if key.ResourceVersion == req.ResourceVersion {
|
||||
exactKeys = append(exactKeys, key)
|
||||
@@ -528,7 +493,7 @@ func filterHistoryKeysByVersion(historyKeys []DataKey, req *resourcepb.ListReque
|
||||
return exactKeys, nil
|
||||
case resourcepb.ResourceVersionMatchV2_NotOlderThan:
|
||||
if req.ResourceVersion > 0 {
|
||||
var filteredKeys []DataKey
|
||||
filteredKeys := make([]DataKey, 0, len(historyKeys))
|
||||
for _, key := range historyKeys {
|
||||
if key.ResourceVersion >= req.ResourceVersion {
|
||||
filteredKeys = append(filteredKeys, key)
|
||||
@@ -538,7 +503,7 @@ func filterHistoryKeysByVersion(historyKeys []DataKey, req *resourcepb.ListReque
|
||||
}
|
||||
default:
|
||||
if req.ResourceVersion > 0 {
|
||||
var filteredKeys []DataKey
|
||||
filteredKeys := make([]DataKey, 0, len(historyKeys))
|
||||
for _, key := range historyKeys {
|
||||
if key.ResourceVersion <= req.ResourceVersion {
|
||||
filteredKeys = append(filteredKeys, key)
|
||||
@@ -564,7 +529,7 @@ func applyLiveHistoryFilter(filteredKeys []DataKey, req *resourcepb.ListRequest)
|
||||
}
|
||||
}
|
||||
if latestDeleteRV > 0 {
|
||||
var liveKeys []DataKey
|
||||
liveKeys := make([]DataKey, 0, len(filteredKeys))
|
||||
for _, key := range filteredKeys {
|
||||
if key.ResourceVersion > latestDeleteRV {
|
||||
liveKeys = append(liveKeys, key)
|
||||
@@ -594,7 +559,7 @@ func applyPagination(keys []DataKey, lastSeenRV int64, sortAscending bool) []Dat
|
||||
return keys
|
||||
}
|
||||
|
||||
var pagedKeys []DataKey
|
||||
pagedKeys := make([]DataKey, 0, len(keys))
|
||||
for _, key := range keys {
|
||||
if sortAscending && key.ResourceVersion > lastSeenRV {
|
||||
pagedKeys = append(pagedKeys, key)
|
||||
@@ -666,7 +631,7 @@ func (k *kvStorageBackend) listModifiedSinceDataStore(ctx context.Context, key N
|
||||
return func(yield func(*ModifiedResource, error) bool) {
|
||||
var lastSeenResource *ModifiedResource
|
||||
var lastSeenDataKey DataKey
|
||||
for dataKey, err := range k.dataStore.Keys(ctx, ListRequestKey{Namespace: key.Namespace, Group: key.Group, Resource: key.Resource}) {
|
||||
for dataKey, err := range k.dataStore.Keys(ctx, ListRequestKey{Namespace: key.Namespace, Group: key.Group, Resource: key.Resource}, SortOrderAsc) {
|
||||
if err != nil {
|
||||
yield(&ModifiedResource{}, err)
|
||||
return
|
||||
@@ -767,7 +732,14 @@ func (k *kvStorageBackend) listModifiedSinceEventStore(ctx context.Context, key
|
||||
}
|
||||
seen[evtKey.Name] = struct{}{}
|
||||
|
||||
value, err := k.getValueFromDataStore(ctx, DataKey(evtKey))
|
||||
value, err := k.getValueFromDataStore(ctx, DataKey{
|
||||
Group: evtKey.Group,
|
||||
Resource: evtKey.Resource,
|
||||
Namespace: evtKey.Namespace,
|
||||
Name: evtKey.Name,
|
||||
ResourceVersion: evtKey.ResourceVersion,
|
||||
Action: evtKey.Action,
|
||||
})
|
||||
if err != nil {
|
||||
yield(&ModifiedResource{}, err)
|
||||
return
|
||||
@@ -820,7 +792,7 @@ func (k *kvStorageBackend) ListHistory(ctx context.Context, req *resourcepb.List
|
||||
Group: key.Group,
|
||||
Resource: key.Resource,
|
||||
Name: key.Name,
|
||||
}) {
|
||||
}, SortOrderAsc) {
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
@@ -872,7 +844,7 @@ func (k *kvStorageBackend) ListHistory(ctx context.Context, req *resourcepb.List
|
||||
// processTrashEntries handles the special case of listing deleted items (trash)
|
||||
func (k *kvStorageBackend) processTrashEntries(ctx context.Context, req *resourcepb.ListRequest, fn func(ListIterator) error, historyKeys []DataKey, lastSeenRV int64, sortAscending bool, listRV int64) (int64, error) {
|
||||
// Filter to only deleted entries
|
||||
var deletedKeys []DataKey
|
||||
deletedKeys := make([]DataKey, 0, len(historyKeys))
|
||||
for _, key := range historyKeys {
|
||||
if key.Action == DataActionDeleted {
|
||||
deletedKeys = append(deletedKeys, key)
|
||||
@@ -881,14 +853,14 @@ func (k *kvStorageBackend) processTrashEntries(ctx context.Context, req *resourc
|
||||
|
||||
// Check if the resource currently exists (is live)
|
||||
// If it exists, don't return any trash entries
|
||||
_, err := k.metaStore.GetLatestResourceKey(ctx, MetaGetRequestKey{
|
||||
Namespace: req.Options.Key.Namespace,
|
||||
_, err := k.dataStore.GetLatestResourceKey(ctx, GetRequestKey{
|
||||
Group: req.Options.Key.Group,
|
||||
Resource: req.Options.Key.Resource,
|
||||
Namespace: req.Options.Key.Namespace,
|
||||
Name: req.Options.Key.Name,
|
||||
})
|
||||
|
||||
var trashKeys []DataKey
|
||||
trashKeys := make([]DataKey, 0, 1)
|
||||
if errors.Is(err, ErrNotFound) {
|
||||
// Resource doesn't exist currently, so we can return the latest delete
|
||||
// Find the latest delete event
|
||||
@@ -1045,12 +1017,13 @@ func (k *kvStorageBackend) WatchWriteEvents(ctx context.Context) (<-chan *Writte
|
||||
for event := range notifierEvents {
|
||||
// fetch the data
|
||||
dataReader, err := k.dataStore.Get(ctx, DataKey{
|
||||
Namespace: event.Namespace,
|
||||
Group: event.Group,
|
||||
Resource: event.Resource,
|
||||
Namespace: event.Namespace,
|
||||
Name: event.Name,
|
||||
ResourceVersion: event.ResourceVersion,
|
||||
Action: event.Action,
|
||||
Folder: event.Folder,
|
||||
})
|
||||
if err != nil || dataReader == nil {
|
||||
k.log.Error("failed to get data for event", "error", err)
|
||||
@@ -1092,48 +1065,8 @@ func (k *kvStorageBackend) WatchWriteEvents(ctx context.Context) (<-chan *Writte
|
||||
}
|
||||
|
||||
// GetResourceStats returns resource stats within the storage backend.
|
||||
// TODO: this isn't very efficient, we should use a more efficient algorithm.
|
||||
func (k *kvStorageBackend) GetResourceStats(ctx context.Context, namespace string, minCount int) ([]ResourceStats, error) {
|
||||
stats := make([]ResourceStats, 0)
|
||||
res := make(map[string]map[string]bool)
|
||||
rvs := make(map[string]int64)
|
||||
|
||||
// Use datastore.Keys to get all data keys for the namespace
|
||||
for dataKey, err := range k.dataStore.Keys(ctx, ListRequestKey{Namespace: namespace}) {
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
key := fmt.Sprintf("%s/%s/%s", dataKey.Namespace, dataKey.Group, dataKey.Resource)
|
||||
if _, ok := res[key]; !ok {
|
||||
res[key] = make(map[string]bool)
|
||||
rvs[key] = 1
|
||||
}
|
||||
res[key][dataKey.Name] = dataKey.Action != DataActionDeleted
|
||||
rvs[key] = dataKey.ResourceVersion
|
||||
}
|
||||
|
||||
for key, names := range res {
|
||||
parts := strings.Split(key, "/")
|
||||
count := int64(0)
|
||||
for _, exists := range names {
|
||||
if exists {
|
||||
count++
|
||||
}
|
||||
}
|
||||
if count <= int64(minCount) {
|
||||
continue
|
||||
}
|
||||
stats = append(stats, ResourceStats{
|
||||
NamespacedResource: NamespacedResource{
|
||||
Namespace: parts[0],
|
||||
Group: parts[1],
|
||||
Resource: parts[2],
|
||||
},
|
||||
Count: count,
|
||||
ResourceVersion: rvs[key],
|
||||
})
|
||||
}
|
||||
return stats, nil
|
||||
return k.dataStore.GetResourceStats(ctx, namespace, minCount)
|
||||
}
|
||||
|
||||
// readAndClose reads all data from a ReadCloser and ensures it's closed,
|
||||
|
||||
Reference in New Issue
Block a user