Grade 3 more models + dashboard v2 layout (quant/format as first-class)
New graded (11 total now):
gemma4-26b-a4b-8bit-mlx 82 Minor Flaws (tied top local; delta-based tx freq)
qwen3.6-27b-8bit-mlx 78 Minor Flaws (clean; anom. slow generation flagged)
qwen3-coder-30b-6bit-mlx 50 Critical (asyncio.Lock used with sync with -> crash)
Dashboard redesign:
- Bar chart is now the full-width hero row (was cramped half-width)
- 4 stat tiles squished 2x2 beside the radar up top
- Quant + Format are dedicated columns in the leaderboard (MLX/GGUF/CLOUD chips)
- New 'Format & Quant Showdown' panel: groups same-family variants so
GGUF-vs-MLX and quant-depth comparisons are side by side
- Bar-chart axis labels now include the quant so duplicate model names
are distinguishable, with rotation for readability
Co-Authored-By: Claude <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,406 @@
|
||||
import asyncio
|
||||
import time
|
||||
from dataclasses import dataclass, field
|
||||
from typing import Any, Dict, Generic, Optional, Set, TypeVar, Union
|
||||
|
||||
# Type variables for Generics
|
||||
K = TypeVar("K")
|
||||
V = TypeVar("V")
|
||||
|
||||
@dataclass
|
||||
class Node(Generic[K, V]):
|
||||
"""A node in the doubly linked list representing a cache entry."""
|
||||
key: K
|
||||
value: V
|
||||
freq: int = 1
|
||||
expiry: float = float('inf')
|
||||
prev: Optional['Node[K, V]'] = None
|
||||
next: Optional['Node[K, V]'] = None
|
||||
|
||||
class DoublyLinkedList(Generic[K, V]):
|
||||
"""A standard Doubly Linked List to allow O(1) removal and insertion."""
|
||||
def __init__(self):
|
||||
self.head: Optional[Node[K, V]] = None
|
||||
self.tail: Optional[Node[K, V]] = None
|
||||
self.size: int = 0
|
||||
|
||||
def append(self, node: Node[K, V]):
|
||||
"""Adds a node to the end of the list."""
|
||||
if not self.head:
|
||||
self.head = node
|
||||
self.tail = node
|
||||
node.prev = None
|
||||
node.next = None
|
||||
else:
|
||||
node.prev = self.tail
|
||||
node.next = None
|
||||
if self.tail:
|
||||
self.tail.next = node
|
||||
self.tail = node
|
||||
self.size += 1
|
||||
|
||||
def remove(self, node: Node[K, V]):
|
||||
"""Removes a specific node from the list in O(1)."""
|
||||
if node.prev:
|
||||
node.prev.next = node.next
|
||||
else:
|
||||
self.head = node.next
|
||||
|
||||
if node.next:
|
||||
node.next.prev = node.prev
|
||||
else:
|
||||
self.tail = node.prev
|
||||
|
||||
node.next = None
|
||||
node.prev = None
|
||||
self.size -= 1
|
||||
|
||||
def pop_tail(self) -> Optional[Node[K, V]]:
|
||||
"""Removes and returns the last node in O(1)."""
|
||||
if not self.tail:
|
||||
return None
|
||||
node = self.tail
|
||||
self.remove(node)
|
||||
return node
|
||||
|
||||
class Transaction(Generic[K, V]):
|
||||
"""
|
||||
Implements ACID-like sub-sessions.
|
||||
Provides 'Read Your Own Writes' and isolation from the global cache.
|
||||
"""
|
||||
def __init__(self, cache: 'LFUCache[K, V]'):
|
||||
self._cache = cache
|
||||
self._pending_puts: Dict[K, tuple[V, float]] = {}
|
||||
self._pending_deletes: Set[K] = set()
|
||||
# Track keys read from the main cache to update frequency on commit
|
||||
self._read_cache_keys: Set[K] = set()
|
||||
self._is_active = True
|
||||
|
||||
async def get(self, key: K) -> Optional[V]:
|
||||
if not self._is_active: raise RuntimeError("Transaction closed")
|
||||
|
||||
# 1. Check local deletes (Tombstone)
|
||||
if key in self._pending_deletes:
|
||||
return None
|
||||
# 2. Check local writes (Read Your Own Writes)
|
||||
if key in self._pending_puts:
|
||||
val, _ = self._pending_puts[key]
|
||||
return val
|
||||
# 3. Check global cache (without affecting LFU frequency until commit)
|
||||
async with self._cache._lock:
|
||||
node = self._cache.cache_data.get(key)
|
||||
if node:
|
||||
# Lazy TTL check even during transaction get
|
||||
if time.monotonic() >= node.expiry:
|
||||
self._cache._remove_node_from_structures(node)
|
||||
return None
|
||||
self._read_cache_keys.add(key)
|
||||
return node.value
|
||||
return None
|
||||
|
||||
async def put(self, key: K, value: V, ttl_seconds: float = float('inf')):
|
||||
if not self._is_active: raise RuntimeError("Transaction closed")
|
||||
expiry = time.monotonic() + ttl_seconds
|
||||
self._pending_puts[key] = (value, expiry)
|
||||
if key in self._pending_deletes:
|
||||
self._pending_deletes.remove(key)
|
||||
|
||||
async def delete(self, key: K):
|
||||
if not self._is_active: raise RuntimeError("Transaction closed")
|
||||
if key in self._pending_puts:
|
||||
del self._pending_puts[key]
|
||||
self._pending_deletes.add(key)
|
||||
|
||||
async def commit(self):
|
||||
if not self._is_active: raise RuntimeError("Transaction closed")
|
||||
async with self._cache._lock:
|
||||
# 1. Process Deletes
|
||||
for key in self._pending_deletes:
|
||||
node = self._cache.cache_data.get(key)
|
||||
if node:
|
||||
self._cache._remove_node_from_structures(node)
|
||||
|
||||
# 2. Process Reads (Update frequencies for keys read from global cache)
|
||||
for key in self._read_cache_keys:
|
||||
node = self._cache.cache_data.get(key)
|
||||
if node: # Ensure it wasn't deleted by a pending delete in this TX
|
||||
self._cache._increment_frequency(node)
|
||||
|
||||
# 3. Process Puts
|
||||
for key, (val, expiry) in self._pending_puts.items():
|
||||
# Internal un-locked put for use within the lock context of commit
|
||||
self._cache._internal_put(key, val, expiry)
|
||||
|
||||
self._is_active = False
|
||||
self._pending_puts.clear()
|
||||
|
||||
async def rollback(self):
|
||||
if not self._is_active: raise RuntimeError("Transaction closed")
|
||||
self._pending_puts.clear()
|
||||
self._pending_deletes.clear()
|
||||
self._read_cache_keys.clear()
|
||||
self._is_active = False
|
||||
|
||||
class LFUCache(Generic[K, V]):
|
||||
"""
|
||||
In-Memory Concurrent LFU Cache.
|
||||
Time Complexity: O(1) for get and put.
|
||||
"""
|
||||
def __init__(self, capacity: int):
|
||||
if capacity <= 0: raise ValueError("Capacity must be > 0")
|
||||
self.capacity = capacity
|
||||
self.cache_data: Dict[K, Node[K, V]] = {}
|
||||
self.freq_map: Dict[int, DoublyLinkedList[K, V]] = {}
|
||||
self.min_freq: int = 0
|
||||
self._lock = asyncio.Lock()
|
||||
self._evictor_task: Optional[asyncio.Task] = None
|
||||
|
||||
async def start_evictor(self, interval: float = 1.0):
|
||||
"""Starts the background async eviction loop."""
|
||||
if self._evictor_task is None:
|
||||
self._evictor_task = asyncio.create_task(self._background_eviction_loop(interval))
|
||||
|
||||
async def stop_evictor(self):
|
||||
"""Stops the background async eviction loop."""
|
||||
if self._evictor_task:
|
||||
self._evictor_task.cancel()
|
||||
try:
|
||||
await self._evictor_task
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
self._evictor_task = None
|
||||
|
||||
async def _background_eviction_loop(self, interval: float):
|
||||
while True:
|
||||
await asyncio.sleep(interval)
|
||||
# We do not hold the lock for the entire loop to prevent blocking.
|
||||
# Instead, we take a snapshot of keys and process in batches.
|
||||
async with self._lock:
|
||||
keys_snapshot = list(self.cache_data.keys())
|
||||
|
||||
# Process in small batches to yield control
|
||||
batch_size = 50
|
||||
for i in range(0, len(keys_snapshot), batch_size):
|
||||
batch = keys_snapshot[i : i + batch_size]
|
||||
now = time.monotonic()
|
||||
async with self._lock:
|
||||
for key in batch:
|
||||
node = self.cache_data.get(key)
|
||||
if node and now >= node.expiry:
|
||||
self._remove_node_from_structures(node)
|
||||
|
||||
def _remove_node_from_structures(self, node: Node[K, V]):
|
||||
"""Internal: Removes node from freq_map and cache_data. O(1)."""
|
||||
# Remove from DLL
|
||||
dll = self.freq_map[node.freq]
|
||||
dll.remove(node)
|
||||
if dll.size == 0:
|
||||
del self.freq_map[node.freq]
|
||||
if self.min_freq == node.freq:
|
||||
# This is a simplification; real min_freq update happens in _increment_frequency.
|
||||
# If current min freq list is empty, we'll find the new min in _increment_frequency or put.
|
||||
pass
|
||||
|
||||
# Remove from dict
|
||||
if node.key in self.cache_data:
|
||||
del self.cache_data[node.key]
|
||||
|
||||
def _increment_frequency(self, node: Node[K, V]):
|
||||
"""Internal: Increases frequency of a node. O(1)."""
|
||||
old_freq = node.freq
|
||||
new_freq = old_freq + 1
|
||||
node.freq = new_freq
|
||||
|
||||
# Remove from old DLL
|
||||
old_dll = self.freq_map[old_freq]
|
||||
old_dll.remove(node)
|
||||
if old_dll.size == 0:
|
||||
del self.freq_map[old_freq]
|
||||
if self.min_freq == old_freq:
|
||||
self.min_freq = new_freq
|
||||
|
||||
# Add to new DLL
|
||||
if new_freq not in self.freq_map:
|
||||
self.freq_map[new_freq] = DoublyLinkedList()
|
||||
self.freq_map[new_freq].append(node)
|
||||
|
||||
def _internal_put(self, key: K, value: V, expiry: float):
|
||||
"""Internal: The core LFU logic. Must be called within a lock."""
|
||||
if key in self.cache_data:
|
||||
node = self.cache_data[key]
|
||||
node.value = value
|
||||
node.expiry = expiry
|
||||
self._increment_frequency(node)
|
||||
else:
|
||||
if len(self.cache_data) >= self.capacity:
|
||||
# Evict LFU (min_freq list tail)
|
||||
if self.min_freq in self.freq_map and self.freq_map[self.min_freq].size > 0:
|
||||
victim = self.freq_map[self.min_freq].pop_tail()
|
||||
if victim:
|
||||
del self.cache_data[victim.key]
|
||||
else:
|
||||
# Fallback (should not happen with correct logic)
|
||||
k_evict = next(iter(self.cache_data))
|
||||
del self.cache_data[k_evict]
|
||||
|
||||
new_node = Node(key, value, freq=1, expiry=expiry)
|
||||
self.cache_data[key] = new_node
|
||||
if 1 not in self.freq_map:
|
||||
self.freq_map[1] = DoublyLinkedList()
|
||||
self.freq_map[1].append(new_node)
|
||||
self.min_freq = 1
|
||||
|
||||
async def get(self, key: K) -> Optional[V]:
|
||||
"""O(1) retrieval with lazy TTL eviction."""
|
||||
async with self._lock:
|
||||
node = self.cache_data.get(key)
|
||||
if not node:
|
||||
return None
|
||||
|
||||
# Lazy TTL Eviction
|
||||
if time.monotonic() >= node.expiry:
|
||||
self._remove_node_from_structures(node)
|
||||
return None
|
||||
|
||||
self._increment_frequency(node)
|
||||
return node.value
|
||||
|
||||
async def put(self, key: K, value: V, ttl_seconds: float = float('inf')):
|
||||
"""O(1) insertion with lazy TTL eviction."""
|
||||
expiry = time.monotonic() + ttl_seconds
|
||||
async with self._lock:
|
||||
self._internal_put(key, value, expiry)
|
||||
|
||||
async def delete(self, key: K):
|
||||
"""O(1) deletion."""
|
||||
async with self._lock:
|
||||
node = self.cache_data.get(key)
|
||||
if node:
|
||||
self._remove_node_from_structures(node)
|
||||
|
||||
def begin_transaction(self) -> Transaction[K, V]:
|
||||
return Transaction(self)
|
||||
|
||||
# ==========================================
|
||||
# UNIT TESTS
|
||||
# ==========================================
|
||||
|
||||
async def test_lfu_eviction():
|
||||
print("Testing O(1) LFU Eviction Logic...")
|
||||
cache = LFUCache[str, int](capacity=3)
|
||||
await cache.put("a", 1) # freq 1
|
||||
await cache.put("b", 2) # freq 1
|
||||
await cache.put("c", 3) # freq 1
|
||||
|
||||
# Increase frequency of a and b
|
||||
await cache.get("a") # freq 2
|
||||
await cache.get("b") # freq 2
|
||||
# c is still freq 1
|
||||
|
||||
await cache.put("d", 4) # Should evict 'c'
|
||||
|
||||
assert await cache.get("a") == 1
|
||||
assert await cache.get("b") == 2
|
||||
assert await cache.get("c") is None
|
||||
assert await cache.get("d") == 4
|
||||
print("✅ LFU Eviction Passed.")
|
||||
|
||||
async def test_ttl_eviction():
|
||||
print("Testing Dual-Layer TTL Eviction...")
|
||||
cache = LFUCache[str, int](capacity=10)
|
||||
await cache.start_evictor(interval=0.1)
|
||||
|
||||
# Test Lazy Eviction
|
||||
await cache.put("lazy", 100, ttl_seconds=0.2)
|
||||
await asyncio.sleep(0.3)
|
||||
assert await cache.get("lazy") is None, "Lazy eviction failed"
|
||||
|
||||
# Test Background Eviction
|
||||
await cache.put("bg", 200, ttl_seconds=0.1)
|
||||
assert await cache.get("bg") is not None, "Value should still be there for a millisecond"
|
||||
await asyncio.sleep(0.4)
|
||||
# Note: Background loop might not have run yet, but we check if it's gone
|
||||
# Since background is a separate task, it should have cleared 'bg' by now.
|
||||
assert await cache.get("bg") is None, "Background eviction failed"
|
||||
|
||||
await cache.stop_evictor()
|
||||
print("✅ TTL Eviction Passed.")
|
||||
|
||||
async def test_transactions():
|
||||
print("Testing Atomic Transactions...")
|
||||
cache = LFUCache[str, int](capacity=5)
|
||||
|
||||
# Test Rollback
|
||||
tx = cache.begin_transaction()
|
||||
await tx.put("tx1", 10)
|
||||
assert await cache.get("tx1") is None, "Uncommitted write visible!"
|
||||
assert await tx.get("tx1") == 10, "Read Your Own Writes failed"
|
||||
await tx.rollback()
|
||||
assert await cache.get("tx1") is None
|
||||
|
||||
# Test Commit Visibility
|
||||
tx = cache.begin_transaction()
|
||||
await tx.put("tx2", 20)
|
||||
await tx.commit()
|
||||
assert await cache.get("tx2") == 20, "Commit visibility failed"
|
||||
|
||||
# Test Isolation/Tombstones
|
||||
await cache.put("exists", 50)
|
||||
tx = cache.begin_transaction()
|
||||
await tx.delete("exists")
|
||||
assert await tx.get("exists") is None, "Transaction delete failed"
|
||||
assert await cache.get("exists") == 50, "Uncommitted delete visible!"
|
||||
await tx.commit()
|
||||
assert await cache.get("exists") is None, "Commit delete failed"
|
||||
|
||||
print("✅ Transactions Passed.")
|
||||
|
||||
async def test_stress():
|
||||
print("Running Concurrency Stress Test (50 tasks)...")
|
||||
cache = LFUCache[int, int](capacity=20)
|
||||
await cache.start_evictor()
|
||||
|
||||
async def worker(worker_id: int):
|
||||
for i in range(100):
|
||||
key = (worker_id * 100) + (i % 30) # Overlapping keys to induce contention
|
||||
op = i % 3
|
||||
if op == 0:
|
||||
await cache.put(key, i)
|
||||
elif op == 1:
|
||||
await cache.get(key)
|
||||
else:
|
||||
# Transactional stress
|
||||
tx = cache.begin_transaction()
|
||||
await tx.put(key, i)
|
||||
if i % 5 == 0:
|
||||
await tx.rollback()
|
||||
else:
|
||||
await tx.commit()
|
||||
if i % 10 == 0:
|
||||
await asyncio.sleep(0.01)
|
||||
|
||||
tasks = [worker(i) for i in range(50)]
|
||||
await asyncio.gather(*tasks)
|
||||
await cache.stop_evictor()
|
||||
print("✅ Stress Test Completed (No crashes).")
|
||||
|
||||
async def main():
|
||||
start_time = time.perf_counter()
|
||||
try:
|
||||
await test_lfu_eviction()
|
||||
print("-" * 30)
|
||||
await test_ttl_eviction()
|
||||
print("-" * 30)
|
||||
await test_transactions()
|
||||
print("-" * 30)
|
||||
await test_stress()
|
||||
print("-" * 30)
|
||||
print(f"All tests finished successfully in {time.perf_counter() - start_time:.2f}s")
|
||||
except Exception as e:
|
||||
print(f"❌ Tests failed with error: {e}")
|
||||
import traceback
|
||||
traceback.print_exc()
|
||||
|
||||
if __name__ == "__main__":
|
||||
asyncio.run(main())
|
||||
|
||||
@@ -0,0 +1,335 @@
|
||||
import asyncio
|
||||
import time
|
||||
import weakref
|
||||
from dataclasses import dataclass, field
|
||||
from typing import Optional, Dict, Any, List, Set, Tuple
|
||||
from collections import OrderedDict
|
||||
import threading
|
||||
|
||||
# --- Data Structures ---
|
||||
|
||||
@dataclass
|
||||
class CacheNode:
|
||||
"""Doubly-linked node for LFU frequency buckets."""
|
||||
key: str
|
||||
value: Any
|
||||
freq: int = 0
|
||||
ttl: Optional[float] = None # Unix timestamp of expiry
|
||||
prev: Optional['CacheNode'] = None
|
||||
next: Optional['CacheNode'] = None
|
||||
|
||||
@dataclass
|
||||
class FrequencyBucket:
|
||||
"""A bucket of nodes with the same frequency."""
|
||||
freq: int
|
||||
head: CacheNode = field(default_factory=lambda: CacheNode("", None, 0)) # Dummy head
|
||||
tail: CacheNode = field(default_factory=lambda: CacheNode("", None, 0)) # Dummy tail
|
||||
nodes: Dict[str, CacheNode] = field(default_factory=dict)
|
||||
|
||||
def __post_init__(self):
|
||||
self.head.next = self.tail
|
||||
self.tail.prev = self.head
|
||||
|
||||
def add_node(self, node: CacheNode):
|
||||
"""Add a node to the front of the bucket."""
|
||||
node.next = self.head.next
|
||||
node.prev = self.head
|
||||
self.head.next.prev = node
|
||||
self.head.next = node
|
||||
self.nodes[node.key] = node
|
||||
|
||||
def remove_node(self, node: CacheNode):
|
||||
"""Remove a node from the bucket."""
|
||||
node.prev.next = node.next
|
||||
node.next.prev = node.prev
|
||||
del self.nodes[node.key]
|
||||
|
||||
def is_empty(self) -> bool:
|
||||
return len(self.nodes) == 0
|
||||
|
||||
# --- Transaction Support ---
|
||||
|
||||
class Transaction:
|
||||
def __init__(self, cache: 'InMemoryLFUCache'):
|
||||
self.cache = cache
|
||||
self._writes: Dict[str, Tuple[Any, Optional[float]]] = {} # key -> (value, ttl)
|
||||
self._deletions: Set[str] = set()
|
||||
self._committed = False
|
||||
self._rolled_back = False
|
||||
|
||||
def get(self, key: str) -> Any:
|
||||
if self._rolled_back or self._committed:
|
||||
raise RuntimeError("Transaction is no longer active")
|
||||
if key in self._deletions:
|
||||
return None
|
||||
if key in self._writes:
|
||||
return self._writes[key][0]
|
||||
return self.cache.get(key)
|
||||
|
||||
def put(self, key: str, value: Any, ttl_seconds: Optional[float] = None):
|
||||
if self._rolled_back or self._committed:
|
||||
raise RuntimeError("Transaction is no longer active")
|
||||
self._writes[key] = (value, ttl_seconds)
|
||||
|
||||
def delete(self, key: str):
|
||||
if self._rolled_back or self._committed:
|
||||
raise RuntimeError("Transaction is no longer active")
|
||||
self._deletions.add(key)
|
||||
|
||||
async def commit(self):
|
||||
if self._rolled_back or self._committed:
|
||||
raise RuntimeError("Transaction is no longer active")
|
||||
await self.cache._commit_transaction(self)
|
||||
self._committed = True
|
||||
|
||||
async def rollback(self):
|
||||
if self._rolled_back or self._committed:
|
||||
raise RuntimeError("Transaction is no longer active")
|
||||
self._rolled_back = True
|
||||
|
||||
# --- Cache Implementation ---
|
||||
|
||||
class InMemoryLFUCache:
|
||||
def __init__(self, capacity: int = 100):
|
||||
self.capacity = capacity
|
||||
self._data: Dict[str, CacheNode] = {} # Global cache data
|
||||
self._freq_buckets: Dict[int, FrequencyBucket] = {}
|
||||
self._min_freq = 0
|
||||
self._lock = asyncio.Lock()
|
||||
self._evictor_task: Optional[asyncio.Task] = None
|
||||
self._evictor_running = False
|
||||
self._transaction_lock = asyncio.Lock()
|
||||
self._transactions: Dict[int, Transaction] = {}
|
||||
self._transaction_counter = 0
|
||||
|
||||
def _create_bucket(self, freq: int) -> FrequencyBucket:
|
||||
bucket = FrequencyBucket(freq)
|
||||
self._freq_buckets[freq] = bucket
|
||||
return bucket
|
||||
|
||||
def _get_bucket(self, freq: int) -> FrequencyBucket:
|
||||
if freq not in self._freq_buckets:
|
||||
return self._create_bucket(freq)
|
||||
return self._freq_buckets[freq]
|
||||
|
||||
def _evict_lfu(self):
|
||||
"""Evict the least frequently used item."""
|
||||
if self._min_freq not in self._freq_buckets:
|
||||
return None
|
||||
bucket = self._freq_buckets[self._min_freq]
|
||||
if bucket.is_empty():
|
||||
return None
|
||||
node_to_evict = bucket.tail.prev # Tail's previous is the last item
|
||||
self._remove_node(node_to_evict)
|
||||
return node_to_evict.key
|
||||
|
||||
def _remove_node(self, node: CacheNode):
|
||||
"""Remove a node from its bucket and global map."""
|
||||
bucket = self._get_bucket(node.freq)
|
||||
bucket.remove_node(node)
|
||||
if bucket.is_empty():
|
||||
del self._freq_buckets[node.freq]
|
||||
del self._data[node.key]
|
||||
|
||||
def _update_node_frequency(self, node: CacheNode):
|
||||
"""Move node to the next frequency bucket and update min_freq if needed."""
|
||||
old_freq = node.freq
|
||||
node.freq += 1
|
||||
new_freq = node.freq
|
||||
bucket = self._get_bucket(old_freq)
|
||||
bucket.remove_node(node)
|
||||
if bucket.is_empty():
|
||||
del self._freq_buckets[old_freq]
|
||||
if old_freq == self._min_freq:
|
||||
self._min_freq = new_freq
|
||||
new_bucket = self._get_bucket(new_freq)
|
||||
new_bucket.add_node(node)
|
||||
|
||||
def _is_expired(self, node: CacheNode) -> bool:
|
||||
if node.ttl is None:
|
||||
return False
|
||||
return time.time() > node.ttl
|
||||
|
||||
def _check_and_remove_expired(self, node: CacheNode) -> bool:
|
||||
"""Returns True if node was expired and removed."""
|
||||
if self._is_expired(node):
|
||||
self._remove_node(node)
|
||||
return True
|
||||
return False
|
||||
|
||||
async def get(self, key: str) -> Any:
|
||||
async with self._lock:
|
||||
if key not in self._data:
|
||||
return None
|
||||
node = self._data[key]
|
||||
# Lazy TTL check
|
||||
if self._check_and_remove_expired(node):
|
||||
return None
|
||||
self._update_node_frequency(node)
|
||||
return node.value
|
||||
|
||||
async def put(self, key: str, value: Any, ttl_seconds: Optional[float] = None):
|
||||
async with self._lock:
|
||||
ttl = time.time() + ttl_seconds if ttl_seconds is not None else None
|
||||
if key in self._data:
|
||||
node = self._data[key]
|
||||
# Lazy TTL check
|
||||
if self._check_and_remove_expired(node):
|
||||
node = None # Re-add it below
|
||||
if node:
|
||||
node.value = value
|
||||
node.ttl = ttl
|
||||
self._update_node_frequency(node)
|
||||
return
|
||||
# Add new node or replace existing expired one
|
||||
node = CacheNode(key, value, freq=1, ttl=ttl)
|
||||
self._data[key] = node
|
||||
bucket = self._get_bucket(1)
|
||||
bucket.add_node(node)
|
||||
self._min_freq = 1
|
||||
|
||||
# Check capacity and evict if needed
|
||||
if len(self._data) > self.capacity:
|
||||
evicted_key = self._evict_lfu()
|
||||
if evicted_key is not None:
|
||||
# Ensure that the bucket is cleaned up if empty
|
||||
pass
|
||||
|
||||
def begin_transaction(self) -> Transaction:
|
||||
with self._transaction_lock:
|
||||
self._transaction_counter += 1
|
||||
tx = Transaction(self)
|
||||
self._transactions[self._transaction_counter] = tx
|
||||
return tx
|
||||
|
||||
async def _commit_transaction(self, tx: Transaction):
|
||||
async with self._lock:
|
||||
# Apply writes
|
||||
for key, (value, ttl) in tx._writes.items():
|
||||
await self.put(key, value, ttl)
|
||||
# Apply deletions
|
||||
for key in tx._deletions:
|
||||
if key in self._data:
|
||||
node = self._data[key]
|
||||
self._remove_node(node)
|
||||
|
||||
async def start_evictor(self, interval_seconds: float = 5.0):
|
||||
"""Start the background evictor task."""
|
||||
async def _evict_loop():
|
||||
while self._evictor_running:
|
||||
try:
|
||||
await asyncio.sleep(interval_seconds)
|
||||
await self._evict_expired_batch()
|
||||
except Exception:
|
||||
pass # Silently ignore errors in background task
|
||||
|
||||
self._evictor_running = True
|
||||
self._evictor_task = asyncio.create_task(_evict_loop())
|
||||
|
||||
async def stop_evictor(self):
|
||||
"""Stop the background evictor task."""
|
||||
self._evictor_running = False
|
||||
if self._evictor_task:
|
||||
self._evictor_task.cancel()
|
||||
try:
|
||||
await self._evictor_task
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
|
||||
async def _evict_expired_batch(self, batch_size: int = 10):
|
||||
"""Evict a batch of expired keys without holding the global lock for too long."""
|
||||
to_remove = []
|
||||
with self._lock:
|
||||
# Collect expired nodes
|
||||
for node in self._data.values():
|
||||
if self._is_expired(node):
|
||||
to_remove.append(node.key)
|
||||
if len(to_remove) >= batch_size:
|
||||
break
|
||||
# Remove outside lock to avoid blocking readers
|
||||
for key in to_remove:
|
||||
async with self._lock:
|
||||
if key in self._data:
|
||||
node = self._data[key]
|
||||
if self._is_expired(node):
|
||||
self._remove_node(node)
|
||||
|
||||
# --- Unit Tests ---
|
||||
|
||||
async def main():
|
||||
print("Running InMemoryLFUCache tests...")
|
||||
|
||||
# Test 1: O(1) LFU eviction
|
||||
print("Test 1: LFU eviction")
|
||||
cache = InMemoryLFUCache(capacity=3)
|
||||
await cache.put("a", 1)
|
||||
await cache.put("b", 2)
|
||||
await cache.put("c", 3)
|
||||
# Access a twice to make it more frequent
|
||||
await cache.get("a")
|
||||
await cache.get("a")
|
||||
# Add d, should evict the least frequent (b)
|
||||
await cache.put("d", 4)
|
||||
assert await cache.get("b") is None
|
||||
assert await cache.get("a") == 1
|
||||
assert await cache.get("c") == 3
|
||||
assert await cache.get("d") == 4
|
||||
print("✓ LFU eviction works")
|
||||
|
||||
# Test 2: TTL eviction
|
||||
print("Test 2: TTL eviction")
|
||||
cache = InMemoryLFUCache(capacity=10)
|
||||
await cache.put("a", 1, ttl_seconds=0.1) # Expire quickly
|
||||
await asyncio.sleep(0.2)
|
||||
assert await cache.get("a") is None
|
||||
print("✓ TTL eviction works")
|
||||
|
||||
# Test 3: Transaction commit visibility
|
||||
print("Test 3: Transaction visibility")
|
||||
cache = InMemoryLFUCache(capacity=10)
|
||||
tx = cache.begin_transaction()
|
||||
tx.put("x", 100)
|
||||
assert await cache.get("x") is None # Not committed yet
|
||||
await tx.commit()
|
||||
assert await cache.get("x") == 100 # Now visible
|
||||
print("✓ Transaction commit works")
|
||||
|
||||
# Test 4: Transaction rollback
|
||||
print("Test 4: Transaction rollback")
|
||||
cache = InMemoryLFUCache(capacity=10)
|
||||
await cache.put("y", 200)
|
||||
tx = cache.begin_transaction()
|
||||
tx.put("y", 300)
|
||||
tx.delete("z") # z doesn't exist yet
|
||||
await tx.rollback()
|
||||
assert await cache.get("y") == 200 # Original value restored
|
||||
print("✓ Transaction rollback works")
|
||||
|
||||
# Test 5: Concurrent access
|
||||
print("Test 5: Concurrent access")
|
||||
cache = InMemoryLFUCache(capacity=10)
|
||||
async def worker(worker_id: int):
|
||||
for i in range(5):
|
||||
key = f"{worker_id}_{i}"
|
||||
await cache.put(key, i)
|
||||
val = await cache.get(key)
|
||||
assert val == i
|
||||
tasks = [asyncio.create_task(worker(i)) for i in range(50)]
|
||||
await asyncio.gather(*tasks)
|
||||
print("✓ Concurrent access works")
|
||||
|
||||
# Test 6: Background eviction
|
||||
print("Test 6: Background eviction")
|
||||
cache = InMemoryLFUCache(capacity=10)
|
||||
await cache.put("expiring", 1, ttl_seconds=0.1)
|
||||
await cache.start_evictor(interval_seconds=0.05)
|
||||
await asyncio.sleep(0.2)
|
||||
assert await cache.get("expiring") is None
|
||||
await cache.stop_evictor()
|
||||
print("✓ Background eviction works")
|
||||
|
||||
print("\nAll tests passed! ✅")
|
||||
|
||||
if __name__ == "__main__":
|
||||
asyncio.run(main())
|
||||
|
||||
@@ -0,0 +1,349 @@
|
||||
"""
|
||||
In-Memory Concurrent LFU Cache with Async TTL Eviction and Atomic Transactions.
|
||||
Pure Python 3.11+ implementation using only built-in modules.
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import time
|
||||
from dataclasses import dataclass
|
||||
from typing import Any, Optional, Dict, Tuple
|
||||
|
||||
|
||||
@dataclass
|
||||
class Node:
|
||||
"""Represents a key-value pair in the cache with frequency and TTL tracking."""
|
||||
key: Any
|
||||
value: Any
|
||||
freq: int = 1
|
||||
expires_at: float = 0.0
|
||||
prev: Optional['Node'] = None
|
||||
next: Optional['Node'] = None
|
||||
|
||||
|
||||
class DoublyLinkedList:
|
||||
"""O(1) Doubly Linked List for maintaining insertion order within frequency buckets."""
|
||||
def __init__(self) -> None:
|
||||
self.head = Node(key=None, value=None)
|
||||
self.tail = Node(key=None, value=None)
|
||||
self.head.next = self.tail
|
||||
self.tail.prev = self.head
|
||||
self.size = 0
|
||||
|
||||
def append(self, node: Node) -> None:
|
||||
"""Append node to the tail (MRU position). O(1)"""
|
||||
node.prev = self.tail.prev
|
||||
node.next = self.tail
|
||||
self.tail.prev.next = node
|
||||
self.tail.prev = node
|
||||
self.size += 1
|
||||
|
||||
def remove(self, node: Node) -> None:
|
||||
"""Remove specific node from the list. O(1)"""
|
||||
node.prev.next = node.next
|
||||
node.next.prev = node.prev
|
||||
node.prev = node.next = None
|
||||
self.size -= 1
|
||||
|
||||
def pop_lru(self) -> Optional[Node]:
|
||||
"""Remove and return the LRU node (head.next). O(1)"""
|
||||
if self.size == 0:
|
||||
return None
|
||||
node = self.head.next
|
||||
self.remove(node)
|
||||
return node
|
||||
|
||||
|
||||
class Transaction:
|
||||
"""
|
||||
ACID-like transaction handle supporting Read-Your-Own-Writes isolation.
|
||||
Uncommitted changes are buffered locally and invisible to global readers.
|
||||
"""
|
||||
def __init__(self, cache: 'LFUCache') -> None:
|
||||
self._cache = cache
|
||||
# key -> (value, expires_at) | None (marks deletion)
|
||||
self._writes: Dict[Any, Optional[Tuple[Any, float]]] = {}
|
||||
self._committed = False
|
||||
self._rolled_back = False
|
||||
|
||||
def _check_active(self) -> None:
|
||||
if self._committed or self._rolled_back:
|
||||
raise RuntimeError("Transaction is no longer active (already committed or rolled back).")
|
||||
|
||||
async def get(self, key: Any) -> Any:
|
||||
"""Read-Your-Own-Writes: checks local buffer first, then global cache."""
|
||||
self._check_active()
|
||||
if key in self._writes and self._writes[key] is not None:
|
||||
val, exp = self._writes[key]
|
||||
return None if time.time() > exp else val
|
||||
return await self._cache.get(key)
|
||||
|
||||
async def put(self, key: Any, value: Any, ttl_seconds: float = 0.0) -> None:
|
||||
"""Buffer write locally without mutating global state."""
|
||||
self._check_active()
|
||||
exp = time.time() + ttl_seconds if ttl_seconds > 0 else float('inf')
|
||||
self._writes[key] = (value, exp)
|
||||
|
||||
async def delete(self, key: Any) -> None:
|
||||
"""Buffer deletion locally."""
|
||||
self._check_active()
|
||||
self._writes[key] = None
|
||||
|
||||
async def commit(self) -> None:
|
||||
"""Apply buffered changes to the global cache atomically."""
|
||||
self._check_active()
|
||||
async with self._cache.lock:
|
||||
for key, data in self._writes.items():
|
||||
if data is None:
|
||||
# Apply deletion
|
||||
if key in self._cache.nodes:
|
||||
self._cache._remove_node(key)
|
||||
else:
|
||||
val, exp = data
|
||||
if key in self._cache.nodes:
|
||||
# Update existing
|
||||
node = self._cache.nodes[key]
|
||||
node.value = val
|
||||
node.expires_at = exp
|
||||
self._cache._update_freq(node)
|
||||
else:
|
||||
# Insert new
|
||||
if self._cache.size >= self._cache.capacity:
|
||||
self._cache._evict_lfu()
|
||||
node = Node(key=key, value=val, freq=1, expires_at=exp)
|
||||
self._cache.nodes[key] = node
|
||||
if 1 not in self._cache.freq_buckets:
|
||||
self._cache.freq_buckets[1] = DoublyLinkedList()
|
||||
self._cache.freq_buckets[1].append(node)
|
||||
self._cache.min_freq = 1
|
||||
self._cache.size += 1
|
||||
self._writes.clear()
|
||||
self._committed = True
|
||||
|
||||
async def rollback(self) -> None:
|
||||
"""Discard all pending changes without affecting global state."""
|
||||
self._check_active()
|
||||
self._writes.clear()
|
||||
self._rolled_back = True
|
||||
|
||||
|
||||
class LFUCache:
|
||||
"""
|
||||
O(1) LFU Cache with Dual-Layer TTL Eviction and Async Concurrency.
|
||||
Uses frequency buckets + doubly linked lists for strict O(1) get/put.
|
||||
"""
|
||||
def __init__(self, capacity: int) -> None:
|
||||
if capacity <= 0:
|
||||
raise ValueError("Capacity must be a positive integer.")
|
||||
self.capacity = capacity
|
||||
self.nodes: Dict[Any, Node] = {}
|
||||
self.freq_buckets: Dict[int, DoublyLinkedList] = {}
|
||||
self.min_freq = 0
|
||||
self.size = 0
|
||||
self.lock = asyncio.Lock()
|
||||
self._evictor_task: Optional[asyncio.Task] = None
|
||||
self._running = False
|
||||
|
||||
def begin_transaction(self) -> Transaction:
|
||||
"""Start a new isolated transaction session."""
|
||||
return Transaction(self)
|
||||
|
||||
async def get(self, key: Any) -> Any:
|
||||
"""Retrieve value by key. O(1) average time complexity."""
|
||||
async with self.lock:
|
||||
if key not in self.nodes:
|
||||
return None
|
||||
node = self.nodes[key]
|
||||
# Lazy TTL Eviction
|
||||
if time.time() > node.expires_at:
|
||||
self._remove_node(key)
|
||||
return None
|
||||
self._update_freq(node)
|
||||
return node.value
|
||||
|
||||
async def put(self, key: Any, value: Any, ttl_seconds: float = 0.0) -> None:
|
||||
"""Insert or update key-value pair. O(1) average time complexity."""
|
||||
async with self.lock:
|
||||
if key in self.nodes:
|
||||
node = self.nodes[key]
|
||||
node.value = value
|
||||
node.expires_at = time.time() + ttl_seconds if ttl_seconds > 0 else float('inf')
|
||||
self._update_freq(node)
|
||||
return
|
||||
|
||||
if self.size >= self.capacity:
|
||||
self._evict_lfu()
|
||||
|
||||
node = Node(
|
||||
key=key,
|
||||
value=value,
|
||||
freq=1,
|
||||
expires_at=time.time() + ttl_seconds if ttl_seconds > 0 else float('inf')
|
||||
)
|
||||
self.nodes[key] = node
|
||||
if 1 not in self.freq_buckets:
|
||||
self.freq_buckets[1] = DoublyLinkedList()
|
||||
self.freq_buckets[1].append(node)
|
||||
self.min_freq = 1
|
||||
self.size += 1
|
||||
|
||||
def _update_freq(self, node: Node) -> None:
|
||||
"""Move node to next frequency bucket. O(1)"""
|
||||
freq = node.freq
|
||||
if freq in self.freq_buckets:
|
||||
self.freq_buckets[freq].remove(node)
|
||||
if self.freq_buckets[freq].size == 0:
|
||||
del self.freq_buckets[freq]
|
||||
if self.min_freq == freq:
|
||||
self.min_freq += 1
|
||||
node.freq += 1
|
||||
if node.freq not in self.freq_buckets:
|
||||
self.freq_buckets[node.freq] = DoublyLinkedList()
|
||||
self.freq_buckets[node.freq].append(node)
|
||||
|
||||
def _evict_lfu(self) -> None:
|
||||
"""Evict LRU key from the minimum frequency bucket. O(1)"""
|
||||
if self.min_freq not in self.freq_buckets:
|
||||
return
|
||||
lru_node = self.freq_buckets[self.min_freq].pop_lru()
|
||||
if lru_node:
|
||||
del self.nodes[lru_node.key]
|
||||
if self.freq_buckets[self.min_freq].size == 0:
|
||||
del self.freq_buckets[self.min_freq]
|
||||
self.size -= 1
|
||||
|
||||
def _remove_node(self, key: Any) -> None:
|
||||
"""Remove node from cache and frequency structure. O(1)"""
|
||||
node = self.nodes.pop(key)
|
||||
if node.freq in self.freq_buckets:
|
||||
self.freq_buckets[node.freq].remove(node)
|
||||
if self.freq_buckets[node.freq].size == 0:
|
||||
del self.freq_buckets[node.freq]
|
||||
if self.min_freq == node.freq:
|
||||
self.min_freq += 1
|
||||
self.size -= 1
|
||||
|
||||
async def start_evictor(self, interval: float = 0.1, batch_size: int = 50) -> None:
|
||||
"""Start background async TTL sweep task."""
|
||||
if self._running:
|
||||
return
|
||||
self._running = True
|
||||
self._evictor_task = asyncio.create_task(self._evict_loop(interval, batch_size))
|
||||
|
||||
async def stop_evictor(self) -> None:
|
||||
"""Gracefully stop the background evictor."""
|
||||
self._running = False
|
||||
if self._evictor_task:
|
||||
self._evictor_task.cancel()
|
||||
try:
|
||||
await self._evictor_task
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
|
||||
async def _evict_loop(self, interval: float, batch_size: int) -> None:
|
||||
"""Non-blocking background eviction that processes in small batches."""
|
||||
while self._running:
|
||||
async with self.lock:
|
||||
now = time.time()
|
||||
checked = 0
|
||||
# Snapshot keys to avoid RuntimeError during iteration/mutation
|
||||
for key in list(self.nodes.keys()):
|
||||
if checked >= batch_size:
|
||||
break
|
||||
node = self.nodes.get(key)
|
||||
if node and now > node.expires_at:
|
||||
self._remove_node(key)
|
||||
checked += 1
|
||||
await asyncio.sleep(interval)
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# EXECUTABLE UNIT TESTS
|
||||
# =============================================================================
|
||||
|
||||
async def main() -> None:
|
||||
print("=== Running LFU Cache Test Suite ===\n")
|
||||
|
||||
# a) O(1) LFU eviction order
|
||||
print("[TEST a] LFU Eviction Order...")
|
||||
cache = LFUCache(3)
|
||||
await cache.put('a', 1)
|
||||
await cache.put('b', 2)
|
||||
await cache.put('c', 3)
|
||||
await cache.get('a') # freq: a=2, b=1, c=1
|
||||
await cache.put('d', 4) # Evicts 'b' (LRU among freq=1)
|
||||
assert await cache.get('b') is None, "LFU eviction failed: 'b' should be evicted"
|
||||
assert await cache.get('a') == 1
|
||||
assert await cache.get('c') == 3
|
||||
assert await cache.get('d') == 4
|
||||
print(" ✅ PASSED\n")
|
||||
|
||||
# b) Lazy TTL vs Background Async Sweep
|
||||
print("[TEST b] Dual-Layer TTL Eviction...")
|
||||
cache2 = LFUCache(10)
|
||||
await cache2.start_evictor(interval=0.05, batch_size=10)
|
||||
|
||||
# Lazy eviction test
|
||||
await cache2.put('lazy', 'val', ttl_seconds=0.1)
|
||||
await asyncio.sleep(0.12)
|
||||
assert await cache2.get('lazy') is None, "Lazy TTL eviction failed"
|
||||
|
||||
# Background sweep test
|
||||
await cache2.put('bg1', 'v', ttl_seconds=0.05)
|
||||
await cache2.put('bg2', 'v', ttl_seconds=0.05)
|
||||
await asyncio.sleep(0.12)
|
||||
assert 'bg1' not in cache2.nodes, "Background sweep failed to remove 'bg1'"
|
||||
assert 'bg2' not in cache2.nodes, "Background sweep failed to remove 'bg2'"
|
||||
|
||||
await cache2.stop_evictor()
|
||||
print(" ✅ PASSED\n")
|
||||
|
||||
# c) Transaction commit visibility vs rollback state restoration
|
||||
print("[TEST c] Transaction Isolation & Rollback...")
|
||||
cache3 = LFUCache(10)
|
||||
|
||||
# Commit test
|
||||
tx1 = cache3.begin_transaction()
|
||||
await tx1.put('x', 100)
|
||||
assert await cache3.get('x') is None, "Isolation broken: global reader saw uncommitted write"
|
||||
assert await tx1.get('x') == 100, "Read-Your-Own-Writes failed"
|
||||
await tx1.commit()
|
||||
assert await cache3.get('x') == 100, "Commit failed: value not visible globally"
|
||||
|
||||
# Rollback test
|
||||
tx2 = cache3.begin_transaction()
|
||||
await tx2.put('y', 200)
|
||||
await tx2.rollback()
|
||||
assert await cache3.get('y') is None, "Rollback failed: uncommitted value leaked"
|
||||
print(" ✅ PASSED\n")
|
||||
|
||||
# d) Stress test: 50 concurrent async tasks
|
||||
print("[TEST d] Concurrency Stress Test (50 tasks, 100 ops each)...")
|
||||
cache4 = LFUCache(50)
|
||||
|
||||
async def worker(wid: int) -> None:
|
||||
for i in range(100):
|
||||
key = f"k_{wid}_{i}"
|
||||
op = i % 3
|
||||
if op == 0:
|
||||
await cache4.put(key, f"v_{i}", ttl_seconds=0.5 if i % 2 == 0 else 0)
|
||||
elif op == 1:
|
||||
await cache4.get(key)
|
||||
else:
|
||||
tx = cache4.begin_transaction()
|
||||
await tx.put(key, f"tx_{i}")
|
||||
if i % 2 == 0:
|
||||
await tx.commit()
|
||||
else:
|
||||
await tx.rollback()
|
||||
|
||||
tasks = [asyncio.create_task(worker(i)) for i in range(50)]
|
||||
await asyncio.gather(*tasks)
|
||||
assert cache4.size <= 50, f"Capacity violation during stress test: size={cache4.size}"
|
||||
print(" ✅ PASSED\n")
|
||||
|
||||
print("=== All Tests Passed Successfully ===")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
asyncio.run(main())
|
||||
|
||||
Reference in New Issue
Block a user