unified-storage: add ListSinceModified to kv store (#110250)

* implement ListKeysSince in event store

 implement data store version

---------

Co-authored-by: Georges Chaudy <chaudyg@gmail.com>
This commit is contained in:
Will Assis
2025-09-04 12:26:40 -04:00
committed by GitHub
co-authored by Georges Chaudy
parent eed8d189ac
commit ea7c370edd
4 changed files with 566 additions and 31 deletions
+181 -2
View File
@@ -452,8 +452,187 @@ func applyPagination(keys []DataKey, lastSeenRV int64, sortAscending bool) []Dat
}
func (k *kvStorageBackend) ListModifiedSince(ctx context.Context, key NamespacedResource, sinceRv int64) (int64, iter.Seq2[*ModifiedResource, error]) {
return 0, func(yield func(*ModifiedResource, error) bool) {
yield(nil, errors.New("not implemented"))
if !key.Valid() {
return 0, func(yield func(*ModifiedResource, error) bool) {
yield(nil, fmt.Errorf("group, resource, and namespace are required"))
}
}
if sinceRv <= 0 {
return 0, func(yield func(*ModifiedResource, error) bool) {
yield(nil, fmt.Errorf("sinceRv must be greater than 0"))
}
}
// Generate a new resource version for the list
listRV := k.snowflake.Generate().Int64()
// Check if sinceRv is older than 1 hour
sinceRvTimestamp := snowflake.ID(sinceRv).Time()
sinceTime := time.Unix(0, sinceRvTimestamp*int64(time.Millisecond))
sinceRvAge := time.Since(sinceTime)
if sinceRvAge > time.Hour {
k.log.Debug("ListModifiedSince using data store", "sinceRv", sinceRv, "sinceRvAge", sinceRvAge)
return listRV, k.listModifiedSinceDataStore(ctx, key, sinceRv)
}
k.log.Debug("ListModifiedSince using event store", "sinceRv", sinceRv, "sinceRvAge", sinceRvAge)
return listRV, k.listModifiedSinceEventStore(ctx, key, sinceRv)
}
func convertEventType(action DataAction) resourcepb.WatchEvent_Type {
switch action {
case DataActionCreated:
return resourcepb.WatchEvent_ADDED
case DataActionUpdated:
return resourcepb.WatchEvent_MODIFIED
case DataActionDeleted:
return resourcepb.WatchEvent_DELETED
default:
panic(fmt.Sprintf("unknown DataAction: %v", action))
}
}
func (k *kvStorageBackend) getValueFromDataStore(ctx context.Context, dataKey DataKey) ([]byte, error) {
raw, err := k.dataStore.Get(ctx, dataKey)
if err != nil {
return []byte{}, err
}
value, err := io.ReadAll(raw)
if err != nil {
return []byte{}, err
}
return value, nil
}
func (k *kvStorageBackend) listModifiedSinceDataStore(ctx context.Context, key NamespacedResource, sinceRv int64) iter.Seq2[*ModifiedResource, error] {
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}) {
if err != nil {
yield(&ModifiedResource{}, err)
return
}
if dataKey.ResourceVersion < sinceRv {
continue
}
if lastSeenResource == nil {
lastSeenResource = &ModifiedResource{
Key: resourcepb.ResourceKey{
Namespace: dataKey.Namespace,
Group: dataKey.Group,
Resource: dataKey.Resource,
Name: dataKey.Name,
},
ResourceVersion: dataKey.ResourceVersion,
Action: convertEventType(dataKey.Action),
}
lastSeenDataKey = dataKey
}
if lastSeenResource.Key.Name != dataKey.Name {
value, err := k.getValueFromDataStore(ctx, lastSeenDataKey)
if err != nil {
yield(&ModifiedResource{}, err)
return
}
lastSeenResource.Value = value
if !yield(lastSeenResource, nil) {
return
}
}
lastSeenResource = &ModifiedResource{
Key: resourcepb.ResourceKey{
Namespace: dataKey.Namespace,
Group: dataKey.Group,
Resource: dataKey.Resource,
Name: dataKey.Name,
},
ResourceVersion: dataKey.ResourceVersion,
Action: convertEventType(dataKey.Action),
}
lastSeenDataKey = dataKey
}
if lastSeenResource != nil {
value, err := k.getValueFromDataStore(ctx, lastSeenDataKey)
if err != nil {
yield(&ModifiedResource{}, err)
return
}
lastSeenResource.Value = value
yield(lastSeenResource, nil)
}
}
}
func (k *kvStorageBackend) listModifiedSinceEventStore(ctx context.Context, key NamespacedResource, sinceRv int64) iter.Seq2[*ModifiedResource, error] {
return func(yield func(*ModifiedResource, error) bool) {
// store all events ordered by RV for the given tenant here
eventKeys := make([]EventKey, 0)
for evtKeyStr, err := range k.eventStore.ListKeysSince(ctx, sinceRv-defaultLookbackPeriod.Nanoseconds()) {
if err != nil {
yield(&ModifiedResource{}, err)
return
}
evtKey, err := ParseEventKey(evtKeyStr)
if err != nil {
yield(&ModifiedResource{}, err)
return
}
if evtKey.ResourceVersion < sinceRv {
continue
}
if evtKey.Group != key.Group || evtKey.Resource != key.Resource || evtKey.Namespace != key.Namespace {
continue
}
eventKeys = append(eventKeys, evtKey)
}
// we only care about the latest revision of every resource in the list
seen := make(map[string]struct{})
for i := len(eventKeys) - 1; i >= 0; i -= 1 {
evtKey := eventKeys[i]
if _, ok := seen[evtKey.Name]; ok {
continue
}
seen[evtKey.Name] = struct{}{}
value, err := k.getValueFromDataStore(ctx, DataKey(evtKey))
if err != nil {
yield(&ModifiedResource{}, err)
return
}
if !yield(&ModifiedResource{
Key: resourcepb.ResourceKey{
Group: evtKey.Group,
Resource: evtKey.Resource,
Namespace: evtKey.Namespace,
Name: evtKey.Name,
},
Action: convertEventType(evtKey.Action),
ResourceVersion: evtKey.ResourceVersion,
Value: value,
}, nil) {
return
}
}
}
}