Skip to content

Deduplication cache only saves hashes for exact pointer timestamp, causing duplicate events #122

Description

@sandesh-js

Description

Grove's deduplication mechanism fails to prevent duplicate events when multiple events in a batch have different published timestamps. The save_hashes() method only saves hashes for events matching the exact pointer value, discarding hashes for all other timestamps.

Environment

  • Grove version: v2.2.0
  • Connector: okta_system_log (affects all chronological connectors)
  • Cache handler: local_file (affects all cache backends)
  • Log order: CHRONOLOGICAL

Root Cause

File: grove/connectors/__init__.py:887-893

def save_hashes(self):
    """Saves the log entry hashes to cache."""
    serialized = json.dumps(
        list(self._hashes.get(self.pointer, set())), separators=(",", ":")
    )
    self._cache.set(self.cache_key(CACHE_KEY_SEEN), self.operation, serialized)

Problem: save_hashes() only persists hashes for self._hashes[self.pointer], but self._hashes is a dict keyed by each unique timestamp in the batch:

self._hashes = {
    "2025-10-17T18:34:47.448Z": {"hash_abc"},
    "2025-10-17T18:34:47.500Z": {"hash_def"},  # pointer is here
    "2025-10-17T18:34:47.600Z": {"hash_xyz"}
}

Only hash_def is saved to cache. Hashes for other timestamps are lost.

Steps to Reproduce

  1. Configure Okta connector with native filtering:

    {
      "name": "okta-filtered",
      "connector": "okta_system_log",
      "processors": [{"name": "filter", "processor": "filter_entries", "filters": [...]}],
      "outputs": {"filtered": "processed"}
    }
  2. Run Grove to collect events (Okta API returns events with millisecond precision timestamps)

  3. Check deduplication cache:

    cat /cache/deduplication.okta_system_log.*/all.cache
    # Shows: [] or only a subset of hashes
  4. Wait 60s for next collection cycle

  5. Observe: Same events are written to output again (duplicates)

Expected Behavior

All event hashes in a collection batch should be saved to cache and used for deduplication on subsequent runs, regardless of their individual timestamps.

Actual Behavior

  • Only hashes for events with published == self.pointer are saved
  • Events with other timestamps are not deduplicated on next run
  • Results in duplicate events being written to output handlers

Impact

  • High: Affects all connectors using time-based APIs with sub-second precision
  • Causes significant data duplication in downstream systems (Kafka, S3, etc.)
  • Increases storage costs and complicates downstream processing
  • Particularly problematic with Okta's inclusive since filter (since >= cursor returns events at boundary)

Affected Connectors

All CHRONOLOGICAL connectors where the API returns events with:

  • Sub-second timestamp precision (milliseconds)
  • Inclusive time filters (since >= timestamp)

Examples: okta_system_log, github_audit_log, gsuite_activities, etc.

Proposed Fix

Option 1: Save all hashes across all timestamps (simple fix)

def save_hashes(self):
    """Saves the log entry hashes to cache."""
    # Save hashes for ALL timestamps, not just self.pointer
    all_hashes = []
    for timestamp, hashes in self._hashes.items():
        all_hashes.extend(list(hashes))
    
    serialized = json.dumps(all_hashes, separators=(",", ":"))
    self._cache.set(self.cache_key(CACHE_KEY_SEEN), self.operation, serialized)

Option 2: Use a sliding window approach (keyed by pointer, but stores last N timestamps)

Option 3: Switch to UUID-based pointers for APIs that provide unique IDs

Workaround

Implement UUID-based deduplication in downstream consumers:

# Kafka consumer example
seen_uuids = set()
for message in consumer:
    event = json.loads(message.value)
    if event['uuid'] in seen_uuids:
        continue  # Skip duplicate
    seen_uuids.add(event['uuid'])
    process(event)

Additional Context

Cache structure shows the issue:

$ cat /cache/deduplication.okta_system_log.*/all.cache
[]  # Empty despite events being collected

Logs confirm pointer is advancing but no hashes stored:

"Pointer successfully saved to cache": "2025-10-17T18:34:47.448Z"
# No "Saving deduplication hashes" debug log appears

Related Code

  • grove/connectors/__init__.py:627-671 - deduplicate_by_hash() method
  • grove/connectors/__init__.py:887-893 - save_hashes() method
  • grove/connectors/__init__.py:850-877 - hashes property getter
  • grove/connectors/__init__.py:245-259 - _run_chronological() calls save_hashes()

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions