1use std::collections::BTreeMap;
55use std::path::{Path, PathBuf};
56use std::sync::{Arc, Mutex};
57
58pub type PluginStoreHandle = Arc<Mutex<PluginStore>>;
62
63const MAGIC: &[u8; 8] = b"LTSTORE\x00";
66
67const SCHEMA_VERSION: u32 = 1;
72
73const MAX_TOTAL_BYTES: usize = 16 * 1024 * 1024;
82
83const FLUSH_EVERY: usize = 64;
86
87const FLUSH_INTERVAL: std::time::Duration = std::time::Duration::from_secs(1);
93
94static LIVE: Mutex<Vec<std::sync::Weak<Mutex<PluginStore>>>> = Mutex::new(Vec::new());
104
105pub fn register(handle: &PluginStoreHandle) {
107 let mut live = LIVE.lock().unwrap_or_else(|p| p.into_inner());
108 live.retain(|w| w.strong_count() > 0);
109 live.push(Arc::downgrade(handle));
110}
111
112pub fn flush_all() {
115 let live: Vec<PluginStoreHandle> = LIVE
116 .lock()
117 .unwrap_or_else(|p| p.into_inner())
118 .iter()
119 .filter_map(std::sync::Weak::upgrade)
120 .collect();
121 for handle in live {
122 let mut store = handle.lock().unwrap_or_else(|p| p.into_inner());
123 if store.dirty > 0 {
124 store.flush();
125 }
126 }
127}
128
129#[derive(Debug)]
135pub struct PluginStore {
136 path: PathBuf,
138 entries: BTreeMap<String, Vec<u8>>,
142 bytes: usize,
145 generation: u64,
149 dirty: usize,
151 last_flush: Option<std::time::Instant>,
154}
155
156impl PluginStore {
157 pub fn open(dir: &Path) -> Self {
164 let path = dir.join("plugin-store.bin");
165 let (entries, generation) = match std::fs::read(&path) {
166 Ok(bytes) => match decode(&bytes) {
167 Ok(decoded) => decoded,
168 Err(reason) => {
169 tracing::debug!(
170 path = %path.display(),
171 %reason,
172 "plugin store unreadable; starting empty"
173 );
174 (BTreeMap::new(), 0)
175 }
176 },
177 Err(e) if e.kind() == std::io::ErrorKind::NotFound => (BTreeMap::new(), 0),
179 Err(error) => {
180 tracing::debug!(
181 path = %path.display(),
182 %error,
183 "plugin store could not be read; starting empty"
184 );
185 (BTreeMap::new(), 0)
186 }
187 };
188 let bytes = entries.values().map(Vec::len).sum();
189 Self {
190 path,
191 entries,
192 bytes,
193 generation,
194 dirty: 0,
195 last_flush: None,
196 }
197 }
198
199 pub fn put(&mut self, key: &str, value: Vec<u8>) -> Result<(), String> {
207 if value.len() > MAX_TOTAL_BYTES {
208 return Err(format!(
209 "store put refused: {} bytes exceeds the {MAX_TOTAL_BYTES}-byte store cap",
210 value.len()
211 ));
212 }
213 let previous = self.entries.get(key).map(Vec::len).unwrap_or(0);
214 if self.bytes - previous + value.len() > MAX_TOTAL_BYTES {
215 tracing::warn!(
216 cap = MAX_TOTAL_BYTES,
217 held = self.bytes,
218 "plugin store full; discarding it wholesale"
219 );
220 self.entries.clear();
221 self.bytes = 0;
222 }
223 self.bytes = self.bytes - self.entries.get(key).map(Vec::len).unwrap_or(0) + value.len();
224 self.entries.insert(key.to_string(), value);
225 self.mutated();
226 Ok(())
227 }
228
229 pub fn get(&self, key: &str) -> Option<Vec<u8>> {
231 self.entries.get(key).cloned()
232 }
233
234 pub fn delete(&mut self, key: &str) -> Result<(), String> {
239 if let Some(old) = self.entries.remove(key) {
240 self.bytes -= old.len();
241 self.mutated();
242 }
243 Ok(())
244 }
245
246 pub fn keys(&self, prefix: &str) -> Vec<String> {
248 self.entries
251 .range(prefix.to_string()..)
252 .take_while(|(k, _)| k.starts_with(prefix))
253 .map(|(k, _)| k.clone())
254 .collect()
255 }
256
257 pub fn generation(&self) -> u64 {
259 self.generation
260 }
261
262 fn mutated(&mut self) {
264 self.generation += 1;
265 self.dirty += 1;
266 let stale = self
267 .last_flush
268 .is_none_or(|at| at.elapsed() >= FLUSH_INTERVAL);
269 if self.dirty >= FLUSH_EVERY || stale {
270 self.flush();
271 }
272 }
273
274 pub fn flush(&mut self) {
282 self.dirty = 0;
283 self.last_flush = Some(std::time::Instant::now());
284 let bytes = encode(&self.entries, self.generation);
285 if let Some(parent) = self.path.parent() {
286 let _ = std::fs::create_dir_all(parent);
287 }
288 let tmp = self.path.with_extension("bin.tmp");
289 if let Err(error) = std::fs::write(&tmp, &bytes) {
290 tracing::debug!(path = %tmp.display(), %error, "plugin store write failed");
291 return;
292 }
293 if let Err(error) = std::fs::rename(&tmp, &self.path) {
294 tracing::debug!(path = %self.path.display(), %error, "plugin store rename failed");
295 let _ = std::fs::remove_file(&tmp);
296 }
297 }
298}
299
300impl Drop for PluginStore {
301 fn drop(&mut self) {
302 if self.dirty > 0 {
303 self.flush();
304 }
305 }
306}
307
308fn encode(entries: &BTreeMap<String, Vec<u8>>, generation: u64) -> Vec<u8> {
313 let mut out = Vec::with_capacity(64 + entries.values().map(Vec::len).sum::<usize>());
314 out.extend_from_slice(MAGIC);
315 out.extend_from_slice(&SCHEMA_VERSION.to_le_bytes());
316 out.extend_from_slice(&generation.to_le_bytes());
317 out.extend_from_slice(&(entries.len() as u64).to_le_bytes());
318 for (key, value) in entries {
319 out.extend_from_slice(&(key.len() as u32).to_le_bytes());
320 out.extend_from_slice(key.as_bytes());
321 out.extend_from_slice(&(value.len() as u32).to_le_bytes());
322 out.extend_from_slice(value);
323 }
324 out
325}
326
327fn decode(bytes: &[u8]) -> Result<(BTreeMap<String, Vec<u8>>, u64), String> {
331 let mut cursor = Cursor { bytes, at: 0 };
332 if cursor.take(MAGIC.len())? != MAGIC {
333 return Err("not a plugin store (bad magic)".into());
334 }
335 let version = cursor.u32()?;
336 if version != SCHEMA_VERSION {
337 return Err(format!(
338 "schema version {version} (expected {SCHEMA_VERSION})"
339 ));
340 }
341 let generation = cursor.u64()?;
342 let count = cursor.u64()?;
343 let mut entries = BTreeMap::new();
344 for _ in 0..count {
345 let key_len = cursor.u32()? as usize;
346 let key = std::str::from_utf8(cursor.take(key_len)?)
347 .map_err(|_| "key is not UTF-8".to_string())?
348 .to_string();
349 let value_len = cursor.u32()? as usize;
350 let value = cursor.take(value_len)?.to_vec();
351 entries.insert(key, value);
352 }
353 Ok((entries, generation))
354}
355
356struct Cursor<'a> {
359 bytes: &'a [u8],
360 at: usize,
361}
362
363impl<'a> Cursor<'a> {
364 fn take(&mut self, n: usize) -> Result<&'a [u8], String> {
365 let end = self.at.checked_add(n).ok_or("length overflow")?;
366 let slice = self.bytes.get(self.at..end).ok_or("truncated")?;
367 self.at = end;
368 Ok(slice)
369 }
370
371 fn u32(&mut self) -> Result<u32, String> {
372 let b = self.take(4)?;
373 Ok(u32::from_le_bytes([b[0], b[1], b[2], b[3]]))
374 }
375
376 fn u64(&mut self) -> Result<u64, String> {
377 let b = self.take(8)?;
378 Ok(u64::from_le_bytes([
379 b[0], b[1], b[2], b[3], b[4], b[5], b[6], b[7],
380 ]))
381 }
382}
383
384#[cfg(test)]
385mod tests {
386 #![allow(clippy::unwrap_used, clippy::panic)]
387
388 use super::*;
389
390 fn open(dir: &Path) -> PluginStore {
391 PluginStore::open(dir)
392 }
393
394 #[test]
395 fn a_value_round_trips() {
396 let dir = tempfile::tempdir().unwrap();
397 let mut store = open(dir.path());
398 assert!(store.get("nodes").is_none(), "nothing is stored yet");
399 store.put("nodes", vec![1, 2, 3]).unwrap();
400 assert_eq!(store.get("nodes"), Some(vec![1, 2, 3]));
401 }
402
403 #[test]
405 fn values_survive_a_restart() {
406 let dir = tempfile::tempdir().unwrap();
407 {
408 let mut store = open(dir.path());
409 store.put("n/ABC", b"a node".to_vec()).unwrap();
410 store.flush();
411 }
412 let reopened = open(dir.path());
413 assert_eq!(reopened.get("n/ABC"), Some(b"a node".to_vec()));
414 }
415
416 #[test]
420 fn a_first_change_is_written_at_once() {
421 let dir = tempfile::tempdir().unwrap();
422 let mut store = open(dir.path());
423 store.put("projects", b"/src/a".to_vec()).unwrap();
424 std::mem::forget(store);
425 assert_eq!(open(dir.path()).get("projects"), Some(b"/src/a".to_vec()));
426 }
427
428 #[test]
431 fn a_burst_waits_and_flush_all_writes_it() {
432 let dir = tempfile::tempdir().unwrap();
433 let handle: PluginStoreHandle = Arc::new(Mutex::new(open(dir.path())));
434 register(&handle);
435 {
436 let mut store = handle.lock().unwrap();
437 store.put("k", b"first".to_vec()).unwrap();
438 store.put("k", b"second".to_vec()).unwrap();
439 }
440 assert_eq!(
441 open(dir.path()).get("k"),
442 Some(b"first".to_vec()),
443 "the second change, a moment later, is held back"
444 );
445 flush_all();
446 assert_eq!(open(dir.path()).get("k"), Some(b"second".to_vec()));
447 std::mem::forget(handle);
448 }
449
450 #[test]
452 fn dropping_the_store_persists_it() {
453 let dir = tempfile::tempdir().unwrap();
454 {
455 let mut store = open(dir.path());
456 store.put("nodes", b"x".to_vec()).unwrap();
457 }
459 assert_eq!(open(dir.path()).get("nodes"), Some(b"x".to_vec()));
460 }
461
462 #[test]
467 fn generation_moves_on_mutation_and_not_on_a_read() {
468 let dir = tempfile::tempdir().unwrap();
469 let mut store = open(dir.path());
470 let start = store.generation();
471
472 store.put("a", vec![1]).unwrap();
473 let after_put = store.generation();
474 assert!(after_put > start, "a put moves the generation");
475
476 let _ = store.get("a");
477 let _ = store.keys("");
478 assert_eq!(
479 store.generation(),
480 after_put,
481 "reads must not move the generation"
482 );
483
484 store.delete("a").unwrap();
485 assert!(
486 store.generation() > after_put,
487 "a delete moves it too — a retraction is a change readers must see"
488 );
489 }
490
491 #[test]
494 fn deleting_an_absent_key_is_ok_and_moves_nothing() {
495 let dir = tempfile::tempdir().unwrap();
496 let mut store = open(dir.path());
497 let before = store.generation();
498 store.delete("never-stored").unwrap();
499 assert_eq!(store.generation(), before);
500 }
501
502 #[test]
506 fn the_generation_survives_a_restart() {
507 let dir = tempfile::tempdir().unwrap();
508 let expected = {
509 let mut store = open(dir.path());
510 store.put("a", vec![1]).unwrap();
511 store.put("b", vec![2]).unwrap();
512 store.flush();
513 store.generation()
514 };
515 assert_eq!(open(dir.path()).generation(), expected);
516 }
517
518 #[test]
519 fn keys_lists_a_prefix_sorted() {
520 let dir = tempfile::tempdir().unwrap();
521 let mut store = open(dir.path());
522 store.put("n/c", vec![]).unwrap();
523 store.put("n/a", vec![]).unwrap();
524 store.put("b/z", vec![]).unwrap();
525 store.put("nodes", vec![]).unwrap();
526
527 assert_eq!(store.keys("n/"), vec!["n/a", "n/c"]);
528 assert_eq!(store.keys("b/"), vec!["b/z"]);
529 assert_eq!(store.keys("").len(), 4, "an empty prefix lists everything");
530 assert!(store.keys("zzz").is_empty());
531 }
532
533 #[test]
536 fn a_key_that_looks_like_a_path_is_just_a_key() {
537 let dir = tempfile::tempdir().unwrap();
538 {
539 let mut store = open(dir.path());
540 store.put("f/../../etc/passwd", b"opaque".to_vec()).unwrap();
541 store.flush();
542 }
543 assert_eq!(
544 open(dir.path()).get("f/../../etc/passwd"),
545 Some(b"opaque".to_vec()),
546 "it round-trips as data"
547 );
548 let written: Vec<_> = std::fs::read_dir(dir.path())
550 .unwrap()
551 .filter_map(|e| e.ok())
552 .map(|e| e.file_name())
553 .collect();
554 assert_eq!(written.len(), 1, "one file, whatever the keys look like");
555 }
556
557 #[test]
561 fn corrupt_bytes_start_empty_rather_than_failing() {
562 let dir = tempfile::tempdir().unwrap();
563 std::fs::write(dir.path().join("plugin-store.bin"), b"not a store at all").unwrap();
564 let mut store = open(dir.path());
565 assert!(store.get("nodes").is_none());
566 store.put("nodes", vec![9]).unwrap();
568 store.flush();
569 assert_eq!(open(dir.path()).get("nodes"), Some(vec![9]));
570 }
571
572 #[test]
575 fn a_truncated_frame_starts_empty() {
576 let dir = tempfile::tempdir().unwrap();
577 {
578 let mut store = open(dir.path());
579 store.put("a", vec![1, 2, 3, 4]).unwrap();
580 store.put("b", vec![5, 6, 7, 8]).unwrap();
581 store.flush();
582 }
583 let path = dir.path().join("plugin-store.bin");
584 let full = std::fs::read(&path).unwrap();
585 std::fs::write(&path, &full[..full.len() - 3]).unwrap();
586
587 let store = open(dir.path());
588 assert!(
589 store.keys("").is_empty(),
590 "a truncated store is discarded whole, not read up to the tear"
591 );
592 }
593
594 #[test]
596 fn a_schema_mismatch_starts_empty() {
597 let dir = tempfile::tempdir().unwrap();
598 let mut stale = Vec::new();
599 stale.extend_from_slice(MAGIC);
600 stale.extend_from_slice(&(SCHEMA_VERSION + 1).to_le_bytes());
601 stale.extend_from_slice(&7u64.to_le_bytes());
602 stale.extend_from_slice(&0u64.to_le_bytes());
603 std::fs::write(dir.path().join("plugin-store.bin"), &stale).unwrap();
604
605 let store = open(dir.path());
606 assert!(store.keys("").is_empty());
607 assert_eq!(store.generation(), 0, "nor is a future generation trusted");
608 }
609
610 #[test]
613 fn the_size_cap_clears_the_store() {
614 let dir = tempfile::tempdir().unwrap();
615 let mut store = open(dir.path());
616 let chunk = vec![0u8; 4 * 1024 * 1024];
617 for i in 0..4 {
618 store.put(&format!("k{i}"), chunk.clone()).unwrap();
619 }
620 assert_eq!(store.keys("").len(), 4, "16 MiB exactly still fits");
621
622 store.put("overflow", chunk.clone()).unwrap();
623 assert_eq!(
624 store.keys(""),
625 vec!["overflow"],
626 "the cap discards wholesale and keeps the newest write"
627 );
628 }
629
630 #[test]
633 fn a_value_larger_than_the_whole_store_is_refused() {
634 let dir = tempfile::tempdir().unwrap();
635 let mut store = open(dir.path());
636 store.put("keep", vec![1]).unwrap();
637 let err = store
638 .put("huge", vec![0u8; MAX_TOTAL_BYTES + 1])
639 .expect_err("a value bigger than the cap cannot be stored");
640 assert!(err.contains("exceeds"), "{err}");
641 assert_eq!(
642 store.get("keep"),
643 Some(vec![1]),
644 "and the refusal costs nothing that was already there"
645 );
646 }
647
648 #[test]
652 fn overwriting_a_key_does_not_grow_the_accounted_size() {
653 let dir = tempfile::tempdir().unwrap();
654 let mut store = open(dir.path());
655 for _ in 0..8 {
656 store.put("nodes", vec![0u8; 4 * 1024 * 1024]).unwrap();
657 }
658 assert_eq!(
659 store.keys(""),
660 vec!["nodes"],
661 "eight rewrites of one 4 MiB key never trip a 16 MiB cap"
662 );
663 }
664}