File size: 1,848 Bytes
1c4c66b | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 | // 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
}
|