- Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathcache.py
More file actions
Latest commit
207 lines (180 loc) · 6.85 KB
/
Copy pathcache.py
File metadata and controls
207 lines (180 loc) · 6.85 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
from __future__ importannotations
importthreading
importtime
fromcollectionsimportOrderedDict
fromcollections.abcimportCallable
fromdataclassesimportdataclass
fromtypingimportAny
MAX_KEY_LENGTH=512
MAX_TTL_SECONDS=31_536_000
classVersionConflict(RuntimeError):
"""Raised when an optimistic cache update targets a stale version."""
@dataclass(frozen=True)
classCacheEntry:
value: Any
expires_at: float|None
version: int
classSkyCache:
"""Thread-safe bounded single-process TTL/LRU cache with optimistic writes."""
def__init__(
self,
default_ttl_seconds: float|None=None,
max_entries: int=10_000,
*,
clock: Callable[[], float] =time.monotonic,
) ->None:
self._validate_ttl(default_ttl_seconds)
ifnot1<=max_entries<=1_000_000:
raiseValueError("max_entries must be between 1 and 1,000,000")
self._default_ttl=default_ttl_seconds
self._max_entries=max_entries
self._clock=clock
self._entries: OrderedDict[str, CacheEntry] =OrderedDict()
self._lock=threading.RLock()
self._version=0
self._stats= {
"hits": 0,
"misses": 0,
"sets": 0,
"deletes": 0,
"evictions": 0,
"expirations": 0,
"cas_conflicts": 0,
}
@staticmethod
def_validate_key(key: str) ->str:
ifnotisinstance(key, str):
raiseTypeError("cache key must be a string")
ifnotkeyorkey.isspace():
raiseValueError("cache key is required")
iflen(key) >MAX_KEY_LENGTH:
raiseValueError(f"cache key cannot exceed {MAX_KEY_LENGTH} characters")
returnkey
@staticmethod
def_validate_ttl(ttl_seconds: float|None) ->None:
ifttl_secondsisNone:
return
ifnotisinstance(ttl_seconds, (int, float)):
raiseTypeError("ttl_seconds must be numeric")
ifttl_seconds<0orttl_seconds>MAX_TTL_SECONDS:
raiseValueError(
f"ttl_seconds must be between 0 and {MAX_TTL_SECONDS} seconds"
)
def_next_version(self) ->int:
self._version+=1
returnself._version
def_expires_at(self, ttl_seconds: float|None) ->float|None:
ttl=self._default_ttlifttl_secondsisNoneelsettl_seconds
self._validate_ttl(ttl)
returnNoneifttlisNoneelseself._clock() +ttl
def_remove_if_expired(self, key: str, now: float) ->bool:
entry=self._entries.get(key)
ifentryisNoneorentry.expires_atisNoneornow<entry.expires_at:
returnFalse
delself._entries[key]
self._stats["expirations"] +=1
returnTrue
def_evict_if_needed(self) ->None:
whilelen(self._entries) >self._max_entries:
self._entries.popitem(last=False)
self._stats["evictions"] +=1
defset(self, key: str, value: Any, ttl_seconds: float|None=None) ->int:
key=self._validate_key(key)
expires_at=self._expires_at(ttl_seconds)
withself._lock:
version=self._next_version()
self._entries[key] =CacheEntry(value=value, expires_at=expires_at, version=version)
self._entries.move_to_end(key)
self._stats["sets"] +=1
self._evict_if_needed()
returnversion
defset_if_absent(
self,
key: str,
value: Any,
ttl_seconds: float|None=None,
) ->tuple[bool, int]:
key=self._validate_key(key)
withself._lock:
self._remove_if_expired(key, self._clock())
existing=self._entries.get(key)
ifexistingisnotNone:
returnFalse, existing.version
returnTrue, self.set(key, value, ttl_seconds)
defcompare_and_set(
self,
key: str,
value: Any,
expected_version: int,
ttl_seconds: float|None=None,
) ->int:
key=self._validate_key(key)
ifexpected_version<1:
raiseValueError("expected_version must be positive")
withself._lock:
self._remove_if_expired(key, self._clock())
existing=self._entries.get(key)
ifexistingisNoneorexisting.version!=expected_version:
self._stats["cas_conflicts"] +=1
raiseVersionConflict("cache entry version does not match")
returnself.set(key, value, ttl_seconds)
defget_entry(self, key: str) ->CacheEntry|None:
key=self._validate_key(key)
withself._lock:
now=self._clock()
ifself._remove_if_expired(key, now):
self._stats["misses"] +=1
returnNone
entry=self._entries.get(key)
ifentryisNone:
self._stats["misses"] +=1
returnNone
self._entries.move_to_end(key)
self._stats["hits"] +=1
returnentry
defget(self, key: str, default: Any=None) ->Any:
entry=self.get_entry(key)
returndefaultifentryisNoneelseentry.value
defdelete(self, key: str) ->bool:
key=self._validate_key(key)
withself._lock:
removed=self._entries.pop(key, None) isnotNone
ifremoved:
self._stats["deletes"] +=1
returnremoved
defpurge_expired(self) ->int:
withself._lock:
now=self._clock()
expired= [
key
forkey, entryinself._entries.items()
ifentry.expires_atisnotNoneandnow>=entry.expires_at
]
forkeyinexpired:
delself._entries[key]
self._stats["expirations"] +=len(expired)
returnlen(expired)
defclear(self) ->int:
withself._lock:
count=len(self._entries)
self._entries.clear()
self._stats["deletes"] +=count
returncount
defstats(self) ->dict[str, int|float]:
withself._lock:
self.purge_expired()
result: dict[str, int|float] =dict(self._stats)
result["size"] =len(self._entries)
result["max_entries"] =self._max_entries
requests=self._stats["hits"] +self._stats["misses"]
result["hit_rate"] = (
self._stats["hits"] /requestsifrequestselse0.0
)
returnresult
def__len__(self) ->int:
withself._lock:
self.purge_expired()
returnlen(self._entries)
# Backward-compatible alias. The implementation is deliberately single-node;
# the historical repository name must not be interpreted as distributed consensus.
DistributedCache=SkyCache