openmeter / dedupe /memorydedupe /memorydedupe.go
Leon4gr45's picture
Upload folder using huggingface_hub (part 3)
1c4c66b verified
Raw
History Blame Contribute Delete
1.85 kB
// Package memorydedupe implements in-memory event deduplication.
package memorydedupe
import (
"context"
"github.com/cloudevents/sdk-go/v2/event"
lru "github.com/hashicorp/golang-lru/v2"
"github.com/openmeterio/openmeter/openmeter/dedupe"
)
const defaultSize = 1024
// Deduplicator implements in-memory event deduplication.
type Deduplicator struct {
store *lru.Cache[string, any]
}
// NewDeduplicator returns a new {Deduplicator}.
func NewDeduplicator(size int) (*Deduplicator, error) {
if size < 1 {
size = defaultSize
}
store, err := lru.New[string, any](size)
if err != nil {
return nil, err
}
return &Deduplicator{
store: store,
}, nil
}
func (d *Deduplicator) IsUnique(ctx context.Context, namespace string, ev event.Event) (bool, error) {
item := dedupe.Item{
Namespace: namespace,
ID: ev.ID(),
Source: ev.Source(),
}
isContained, _ := d.store.ContainsOrAdd(item.Key(), nil)
return !isContained, nil
}
func (d *Deduplicator) CheckUnique(ctx context.Context, item dedupe.Item) (bool, error) {
isContained := d.store.Contains(item.Key())
return !isContained, nil
}
func (d *Deduplicator) Set(ctx context.Context, items ...dedupe.Item) ([]dedupe.Item, error) {
for _, item := range items {
_ = d.store.Add(item.Key(), nil)
}
return nil, nil
}
func (d *Deduplicator) Close() error {
return nil
}
func (d *Deduplicator) CheckUniqueBatch(ctx context.Context, items []dedupe.Item) (dedupe.CheckUniqueBatchResult, error) {
result := dedupe.CheckUniqueBatchResult{
UniqueItems: make(dedupe.ItemSet, len(items)),
AlreadyProcessedItems: make(dedupe.ItemSet, len(items)),
}
for _, item := range items {
if d.store.Contains(item.Key()) {
result.AlreadyProcessedItems[item] = struct{}{}
} else {
result.UniqueItems[item] = struct{}{}
}
}
return result, nil
}