| 12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970 |
- // Package dedupe implements the 60s-window dedupe with dedupe_count
- // return per SPEC §5.
- //
- // Algorithm:
- // key = dedupe:{source_id}:{dedupe_key}
- // SET key 1 NX EX 60
- // OK -> first arrival, return (count=1, isNew=true)
- // nil -> key existed; INCR; return (count=N, isNew=false)
- //
- // M0 uses plain SET NX + INCR. We can swap to a Lua script for
- // atomicity later if we see races.
- package dedupe
- import (
- "context"
- "errors"
- "fmt"
- "time"
- "github.com/redis/go-redis/v9"
- )
- const (
- DefaultWindow = 60 * time.Second
- keyPrefix = "dedupe:"
- )
- type Deduper struct {
- rdb *redis.Client
- window time.Duration
- }
- func New(rdb *redis.Client, window time.Duration) *Deduper {
- if window <= 0 {
- window = DefaultWindow
- }
- return &Deduper{rdb: rdb, window: window}
- }
- // Check atomically claims a dedupe slot. Returns:
- //
- // (true, 1, nil) – first arrival in the window
- // (false, n, nil) – duplicate; n is the 1-indexed count in this window
- // (false, 0, err) – redis error; caller should fail open
- func (d *Deduper) Check(ctx context.Context, sourceID, dedupeKey string) (isNew bool, count uint32, err error) {
- if dedupeKey == "" {
- // No dedupe key => not dedupable. Caller treats it as a new alert.
- return true, 1, nil
- }
- key := keyPrefix + sourceID + ":" + dedupeKey
- ok, err := d.rdb.SetNX(ctx, key, 1, d.window).Result()
- if err != nil {
- return false, 0, fmt.Errorf("dedupe setnx: %w", err)
- }
- if ok {
- return true, 1, nil
- }
- n, err := d.rdb.Incr(ctx, key).Result()
- if err != nil {
- return false, 0, fmt.Errorf("dedupe incr: %w", err)
- }
- if n < 1 {
- n = 1
- }
- return false, uint32(n), nil
- }
- // ErrInvalidSource is returned for a missing sourceID.
- var ErrInvalidSource = errors.New("dedupe: empty source id")
|