dedupe.go 1.7 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970
  1. // Package dedupe implements the 60s-window dedupe with dedupe_count
  2. // return per SPEC §5.
  3. //
  4. // Algorithm:
  5. // key = dedupe:{source_id}:{dedupe_key}
  6. // SET key 1 NX EX 60
  7. // OK -> first arrival, return (count=1, isNew=true)
  8. // nil -> key existed; INCR; return (count=N, isNew=false)
  9. //
  10. // M0 uses plain SET NX + INCR. We can swap to a Lua script for
  11. // atomicity later if we see races.
  12. package dedupe
  13. import (
  14. "context"
  15. "errors"
  16. "fmt"
  17. "time"
  18. "github.com/redis/go-redis/v9"
  19. )
  20. const (
  21. DefaultWindow = 60 * time.Second
  22. keyPrefix = "dedupe:"
  23. )
  24. type Deduper struct {
  25. rdb *redis.Client
  26. window time.Duration
  27. }
  28. func New(rdb *redis.Client, window time.Duration) *Deduper {
  29. if window <= 0 {
  30. window = DefaultWindow
  31. }
  32. return &Deduper{rdb: rdb, window: window}
  33. }
  34. // Check atomically claims a dedupe slot. Returns:
  35. //
  36. // (true, 1, nil) – first arrival in the window
  37. // (false, n, nil) – duplicate; n is the 1-indexed count in this window
  38. // (false, 0, err) – redis error; caller should fail open
  39. func (d *Deduper) Check(ctx context.Context, sourceID, dedupeKey string) (isNew bool, count uint32, err error) {
  40. if dedupeKey == "" {
  41. // No dedupe key => not dedupable. Caller treats it as a new alert.
  42. return true, 1, nil
  43. }
  44. key := keyPrefix + sourceID + ":" + dedupeKey
  45. ok, err := d.rdb.SetNX(ctx, key, 1, d.window).Result()
  46. if err != nil {
  47. return false, 0, fmt.Errorf("dedupe setnx: %w", err)
  48. }
  49. if ok {
  50. return true, 1, nil
  51. }
  52. n, err := d.rdb.Incr(ctx, key).Result()
  53. if err != nil {
  54. return false, 0, fmt.Errorf("dedupe incr: %w", err)
  55. }
  56. if n < 1 {
  57. n = 1
  58. }
  59. return false, uint32(n), nil
  60. }
  61. // ErrInvalidSource is returned for a missing sourceID.
  62. var ErrInvalidSource = errors.New("dedupe: empty source id")