// 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")