mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-06 03:38:23 +00:00
feat(quartz): add SeenIds — a memory-lean event-id dedup filter
A run-scoped "already seen this id" filter for large, mostly-duplicate id streams (a broad relay walk re-receiving the same event from many relays). Keys on the first 128 bits of the id, sliced straight out of the hex with Hex.readLong (table lookups, no parse, no allocation), in one open-addressed LongArray — ~16 bytes/entry and the 64-char String is never retained, so tens of millions of ids cost ~1 GB instead of a HashSet<String>'s ~6 GB. add() is O(1) and synchronized. Lives in the jvmAndroid source set (uses @Synchronized; a 40M-id walk is a server-side concern). Ports the caller's implementation with the parseUnsignedLong hot path swapped for Hex.readLong (~45-70 ns/op cheaper). Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_015YEbdqCRPkszkGCoi89RMt
This commit is contained in:
@@ -0,0 +1,145 @@
|
||||
/*
|
||||
* Copyright (c) 2025 Vitor Pamplona
|
||||
*
|
||||
* Permission is hereby granted, free of charge, to any person obtaining a copy of
|
||||
* this software and associated documentation files (the "Software"), to deal in
|
||||
* the Software without restriction, including without limitation the rights to use,
|
||||
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
|
||||
* Software, and to permit persons to whom the Software is furnished to do so,
|
||||
* subject to the following conditions:
|
||||
*
|
||||
* The above copyright notice and this permission notice shall be included in all
|
||||
* copies or substantial portions of the Software.
|
||||
*
|
||||
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
|
||||
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
|
||||
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
|
||||
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
|
||||
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
|
||||
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
|
||||
*/
|
||||
package com.vitorpamplona.quartz.utils
|
||||
|
||||
/**
|
||||
* A memory-lean "already seen this event id" filter for large, mostly-duplicate id
|
||||
* streams — e.g. a broad relay walk that re-receives the same widely-mirrored event
|
||||
* from dozens of relays. [add] drops a duplicate the moment it arrives, before any
|
||||
* expensive per-event work (signature verification, a store existence check).
|
||||
*
|
||||
* Event ids are SHA-256 hashes — uniform random 256-bit values — so their first 128
|
||||
* bits (two longs of the 32-byte id) are themselves a perfect hash. Keying on those,
|
||||
* the odds of two distinct ids colliding across tens of millions of events is ~1e-22,
|
||||
* so it never wrongly skips a real event (unlike a Bloom filter). The first 128 bits
|
||||
* are sliced straight out of the hex string with [Hex.readLong] — table lookups and
|
||||
* shifts, no text parsing and no allocation.
|
||||
*
|
||||
* Backed by one open-addressed [LongArray] (two longs per slot, `(0,0)` = empty), so
|
||||
* there are NO per-entry objects and the 64-char id [String] is never retained: tens
|
||||
* of millions of ids cost ~16 bytes each (~1 GB at 40M) instead of the ~6 GB a
|
||||
* `HashSet<String>` of 64-char hex would. [add] is O(1) and `@Synchronized`; the lock
|
||||
* is held for nanoseconds, so concurrent producers contend little.
|
||||
*
|
||||
* Not unbounded-safe on its own: call [reset] between passes (or whenever the working
|
||||
* set should be forgotten) so a long-running process can't grow the table forever.
|
||||
*/
|
||||
class SeenIds(
|
||||
initialSlotsPow2: Int = INITIAL_POW2,
|
||||
) {
|
||||
private var mask = 0
|
||||
private var table = LongArray(0)
|
||||
private var count = 0 // non-zero-key entries held in [table]
|
||||
private var zeroSeen = false // the (0,0) key, tracked apart from the empty sentinel
|
||||
private var resizeAt = 0
|
||||
|
||||
init {
|
||||
allocate(1 shl initialSlotsPow2)
|
||||
}
|
||||
|
||||
private fun allocate(slots: Int) {
|
||||
table = LongArray(slots * 2)
|
||||
mask = slots - 1
|
||||
resizeAt = (slots * LOAD).toInt()
|
||||
count = 0
|
||||
}
|
||||
|
||||
/**
|
||||
* Records [idHex] (a 64-char hex event id); returns true if it is NEW this pass
|
||||
* (the caller should process it), false if already seen (the caller should skip
|
||||
* it). A too-short/malformed id returns true — it flows through and downstream
|
||||
* verification drops it — rather than risk collapsing distinct ids.
|
||||
*/
|
||||
@Synchronized
|
||||
fun add(idHex: String): Boolean {
|
||||
// Slice the first 128 bits straight to two longs via Hex's table-lookup
|
||||
// reader — no hex text parsing, no allocation. A string too short to slice
|
||||
// can't be a real 32-byte id, so let it through (verification drops it).
|
||||
if (idHex.length < 32) return true
|
||||
return addKey(Hex.readLong(idHex, 0), Hex.readLong(idHex, 16))
|
||||
}
|
||||
|
||||
private fun addKey(
|
||||
hi: Long,
|
||||
lo: Long,
|
||||
): Boolean {
|
||||
if (hi == 0L && lo == 0L) {
|
||||
// (0,0) is [table]'s empty sentinel, so this one key is tracked apart.
|
||||
if (zeroSeen) return false
|
||||
zeroSeen = true
|
||||
return true
|
||||
}
|
||||
if (count >= resizeAt) grow()
|
||||
var i = (mix(hi, lo).toInt() and mask)
|
||||
while (true) {
|
||||
val s = i * 2
|
||||
val h = table[s]
|
||||
val l = table[s + 1]
|
||||
if (h == 0L && l == 0L) {
|
||||
table[s] = hi
|
||||
table[s + 1] = lo
|
||||
count++
|
||||
return true
|
||||
}
|
||||
if (h == hi && l == lo) return false
|
||||
i = (i + 1) and mask
|
||||
}
|
||||
}
|
||||
|
||||
private fun grow() {
|
||||
val old = table
|
||||
allocate((mask + 1) shl 1) // resets count; zeroSeen is untouched
|
||||
var j = 0
|
||||
while (j < old.size) {
|
||||
val h = old[j]
|
||||
val l = old[j + 1]
|
||||
if (h != 0L || l != 0L) addKey(h, l)
|
||||
j += 2
|
||||
}
|
||||
}
|
||||
|
||||
@Synchronized
|
||||
fun reset() {
|
||||
allocate(1 shl INITIAL_POW2)
|
||||
zeroSeen = false
|
||||
}
|
||||
|
||||
@Synchronized
|
||||
fun size() = count + if (zeroSeen) 1 else 0
|
||||
|
||||
// Ids are already uniform, but avalanche the two halves so the low bits used for
|
||||
// the slot index don't correlate with any particular byte of the hash.
|
||||
private fun mix(
|
||||
hi: Long,
|
||||
lo: Long,
|
||||
): Long {
|
||||
var h = hi xor (lo * -0x61c8864680b583ebL)
|
||||
h = h xor (h ushr 32)
|
||||
h *= -0x7ee3623a03d3f7d7L
|
||||
h = h xor (h ushr 29)
|
||||
return h
|
||||
}
|
||||
|
||||
companion object {
|
||||
private const val LOAD = 0.7
|
||||
private const val INITIAL_POW2 = 20
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,88 @@
|
||||
/*
|
||||
* Copyright (c) 2025 Vitor Pamplona
|
||||
*
|
||||
* Permission is hereby granted, free of charge, to any person obtaining a copy of
|
||||
* this software and associated documentation files (the "Software"), to deal in
|
||||
* the Software without restriction, including without limitation the rights to use,
|
||||
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
|
||||
* Software, and to permit persons to whom the Software is furnished to do so,
|
||||
* subject to the following conditions:
|
||||
*
|
||||
* The above copyright notice and this permission notice shall be included in all
|
||||
* copies or substantial portions of the Software.
|
||||
*
|
||||
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
|
||||
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
|
||||
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
|
||||
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
|
||||
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
|
||||
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
|
||||
*/
|
||||
package com.vitorpamplona.quartz.utils
|
||||
|
||||
import kotlin.test.Test
|
||||
import kotlin.test.assertEquals
|
||||
import kotlin.test.assertFalse
|
||||
import kotlin.test.assertTrue
|
||||
|
||||
class SeenIdsTest {
|
||||
// Vary the FIRST 128 bits (the keyed part): value in the high 16 hex chars, zero tail.
|
||||
private fun id(i: Int) = i.toLong().toString(16).padStart(16, '0') + "0".repeat(48)
|
||||
|
||||
@Test
|
||||
fun `first sight is new, repeats are skipped`() {
|
||||
val seen = SeenIds()
|
||||
assertTrue(seen.add(id(1)), "first time is new")
|
||||
assertFalse(seen.add(id(1)), "same id is a duplicate")
|
||||
assertFalse(seen.add(id(1)), "still a duplicate")
|
||||
assertTrue(seen.add(id(2)), "a different id is new")
|
||||
assertEquals(2, seen.size())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `reset forgets everything`() {
|
||||
val seen = SeenIds()
|
||||
seen.add(id(7))
|
||||
assertFalse(seen.add(id(7)))
|
||||
seen.reset()
|
||||
assertTrue(seen.add(id(7)), "after reset the id is new again")
|
||||
assertEquals(1, seen.size())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `holds many distinct ids across resizes, with exact dedup`() {
|
||||
// Start tiny so it must grow several times.
|
||||
val seen = SeenIds(initialSlotsPow2 = 4)
|
||||
val n = 50_000
|
||||
repeat(n) { assertTrue(seen.add(id(it)), "id $it should be new") }
|
||||
assertEquals(n, seen.size())
|
||||
// Every one is now a duplicate.
|
||||
repeat(n) { assertFalse(seen.add(id(it)), "id $it should be a duplicate") }
|
||||
assertEquals(n, seen.size(), "duplicates don't grow the set")
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `only the first 128 bits key the id (differ past 32 hex chars still dedups)`() {
|
||||
// Same first 128 bits, different tail -> treated as the same (documented tradeoff, ~1e-22 in practice).
|
||||
val seen = SeenIds()
|
||||
val prefix = "%032x".format(42)
|
||||
assertTrue(seen.add(prefix + "0".repeat(32)))
|
||||
assertFalse(seen.add(prefix + "f".repeat(32)), "same 128-bit prefix collapses")
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `the all-zero 128-bit-prefix key (the empty sentinel) is deduped correctly`() {
|
||||
val seen = SeenIds()
|
||||
assertTrue(seen.add("0".repeat(64)), "all-zero id is new the first time")
|
||||
assertFalse(seen.add("0".repeat(64)), "and a duplicate the second")
|
||||
assertTrue(seen.add(id(5)), "a normal id still works alongside it")
|
||||
assertEquals(2, seen.size())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `a malformed id is let through, not skipped`() {
|
||||
val seen = SeenIds()
|
||||
assertTrue(seen.add("not-hex"), "malformed -> flows to verify")
|
||||
assertTrue(seen.add("short"))
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user