1use crate::chunk::{Manifest, verify_chunk};
2use crate::crypto::{CryptoProvider, DefaultProvider, PrivateIdentity};
3use crate::identity::{Contact, ContactBook, ContactCard, IdentityFile};
4use crate::protocol::{decode_cbor, encode_cbor, verify_envelope};
5use crate::{
6 DEFAULT_STORE_QUOTA, DdError, DropEnvelope, DropState, ErrorCode, ObjectId, Ownership, PeerId,
7 Result, STORAGE_SCHEMA_VERSION, TrustState, hex_encode,
8};
9use rusqlite::{Connection, OptionalExtension, params};
10use std::fs;
11use std::path::{Path, PathBuf};
12use std::sync::Mutex;
13use std::time::{SystemTime, UNIX_EPOCH};
14
15#[derive(Debug, Clone)]
16pub struct StorageQuotas {
17 pub maximum: u64,
18 pub reserved_local: u64,
19 pub relay_budget: u64,
20 pub temporary: u64,
21}
22
23impl Default for StorageQuotas {
24 fn default() -> Self {
25 Self {
26 maximum: DEFAULT_STORE_QUOTA,
27 reserved_local: 5 * 1024 * 1024 * 1024,
28 relay_budget: 10 * 1024 * 1024 * 1024,
29 temporary: 5 * 1024 * 1024 * 1024,
30 }
31 }
32}
33
34pub struct Store {
35 root: PathBuf,
36 quotas: StorageQuotas,
37 db: Mutex<Connection>,
38}
39
40impl Store {
41 pub fn open(root: impl AsRef<Path>, quotas: StorageQuotas) -> Result<Self> {
42 let root = root.as_ref().to_path_buf();
43 fs::create_dir_all(root.join("chunks"))?;
44 fs::create_dir_all(root.join("objects"))?;
45 fs::create_dir_all(root.join("manifests"))?;
46 fs::create_dir_all(root.join("identities"))?;
47 let db_path = root.join("meta.sqlite");
48 let db = Connection::open(&db_path).map_err(sql_err)?;
49 db.execute_batch(
50 "
51 PRAGMA journal_mode=WAL;
52 PRAGMA foreign_keys=ON;
53 CREATE TABLE IF NOT EXISTS meta (k TEXT PRIMARY KEY, v TEXT NOT NULL);
54 CREATE TABLE IF NOT EXISTS objects (
55 object_id BLOB PRIMARY KEY,
56 envelope BLOB NOT NULL,
57 source BLOB NOT NULL,
58 created_at INTEGER NOT NULL,
59 expires_at INTEGER NOT NULL,
60 priority INTEGER NOT NULL,
61 hop_count INTEGER NOT NULL,
62 hop_limit INTEGER NOT NULL,
63 state TEXT NOT NULL,
64 ownership TEXT NOT NULL,
65 size INTEGER NOT NULL,
66 replication INTEGER NOT NULL,
67 application TEXT NOT NULL
68 );
69 CREATE TABLE IF NOT EXISTS chunks (
70 content_id BLOB PRIMARY KEY,
71 size INTEGER NOT NULL,
72 refcount INTEGER NOT NULL
73 );
74 CREATE TABLE IF NOT EXISTS object_chunks (
75 object_id BLOB NOT NULL,
76 idx INTEGER NOT NULL,
77 content_id BLOB NOT NULL,
78 present INTEGER NOT NULL,
79 PRIMARY KEY (object_id, idx)
80 );
81 CREATE TABLE IF NOT EXISTS manifests (
82 manifest_id BLOB PRIMARY KEY,
83 object_id BLOB NOT NULL,
84 body BLOB NOT NULL
85 );
86 CREATE TABLE IF NOT EXISTS contacts (
87 peer_id BLOB PRIMARY KEY,
88 card TEXT NOT NULL,
89 trust TEXT NOT NULL
90 );
91 CREATE TABLE IF NOT EXISTS encounters (
92 peer_id BLOB PRIMARY KEY,
93 first_seen INTEGER NOT NULL,
94 last_seen INTEGER NOT NULL,
95 encounter_count INTEGER NOT NULL,
96 bytes_sent INTEGER NOT NULL,
97 bytes_received INTEGER NOT NULL,
98 successful_forwards INTEGER NOT NULL,
99 failed_forwards INTEGER NOT NULL
100 );
101 CREATE TABLE IF NOT EXISTS receipts (
102 receipt_id BLOB PRIMARY KEY,
103 object_id BLOB NOT NULL,
104 kind TEXT NOT NULL,
105 body BLOB NOT NULL
106 );
107 CREATE TABLE IF NOT EXISTS traces (
108 id INTEGER PRIMARY KEY AUTOINCREMENT,
109 object_id BLOB NOT NULL,
110 ts INTEGER NOT NULL,
111 event TEXT NOT NULL
112 );
113 CREATE TABLE IF NOT EXISTS pins (
114 object_id BLOB PRIMARY KEY
115 );
116 CREATE TABLE IF NOT EXISTS tags (
117 object_id BLOB NOT NULL,
118 tag TEXT NOT NULL,
119 PRIMARY KEY (object_id, tag)
120 );
121 CREATE TABLE IF NOT EXISTS aliases (
122 peer_id BLOB PRIMARY KEY,
123 alias TEXT NOT NULL
124 );
125 CREATE TABLE IF NOT EXISTS spaces (
126 name TEXT PRIMARY KEY,
127 body TEXT NOT NULL
128 );
129 CREATE TABLE IF NOT EXISTS channels (
130 name TEXT PRIMARY KEY,
131 filter TEXT NOT NULL,
132 policy TEXT NOT NULL
133 );
134 CREATE TABLE IF NOT EXISTS circles (
135 name TEXT PRIMARY KEY,
136 members TEXT NOT NULL
137 );
138 CREATE TABLE IF NOT EXISTS kv (
139 k TEXT PRIMARY KEY,
140 v TEXT NOT NULL
141 );
142 CREATE INDEX IF NOT EXISTS idx_objects_exp ON objects(expires_at);
143 ",
144 )
145 .map_err(sql_err)?;
146 db.execute(
147 "INSERT OR IGNORE INTO meta(k,v) VALUES('schema', ?1)",
148 params![STORAGE_SCHEMA_VERSION.to_string()],
149 )
150 .map_err(sql_err)?;
151 Ok(Self {
152 root,
153 quotas,
154 db: Mutex::new(db),
155 })
156 }
157
158 pub fn root(&self) -> &Path {
159 &self.root
160 }
161
162 pub fn identity_path(&self) -> PathBuf {
163 self.root.join("identities").join("identity.json")
164 }
165
166 pub fn save_identity(&self, id: &PrivateIdentity) -> Result<()> {
167 let file = IdentityFile::from_private(id);
168 fs::write(
169 self.identity_path(),
170 serde_json::to_vec_pretty(&file).map_err(|e| DdError::crypto(e.to_string()))?,
171 )?;
172 Ok(())
173 }
174
175 pub fn load_identity(&self) -> Result<PrivateIdentity> {
176 let bytes = fs::read(self.identity_path())?;
177 let file: IdentityFile =
178 serde_json::from_slice(&bytes).map_err(|e| DdError::crypto(e.to_string()))?;
179 file.into_private()
180 }
181
182 pub fn init_identity(&self, force: bool) -> Result<PrivateIdentity> {
183 if self.identity_path().exists() && !force {
184 return Err(DdError::protocol(
185 ErrorCode::Ddx0000Internal,
186 "identity exists (use rotate, not silent replace)",
187 ));
188 }
189 let id = PrivateIdentity::generate();
190 self.save_identity(&id)?;
191 Ok(id)
192 }
193
194 pub fn put_contact(&self, card: ContactCard, trust: TrustState) -> Result<PeerId> {
195 let id = card.peer_id()?;
196 let json = card.to_ddcontact()?;
197 self.db
198 .lock()
199 .expect("db")
200 .execute(
201 "INSERT OR REPLACE INTO contacts(peer_id, card, trust) VALUES(?1,?2,?3)",
202 params![
203 id.as_bytes().as_slice(),
204 json,
205 format!("{trust:?}").to_lowercase()
206 ],
207 )
208 .map_err(sql_err)?;
209 Ok(id)
210 }
211
212 pub fn load_contacts(&self) -> Result<ContactBook> {
213 let db = self.db.lock().expect("db");
214 let mut stmt = db
215 .prepare("SELECT card, trust FROM contacts")
216 .map_err(sql_err)?;
217 let rows = stmt
218 .query_map([], |row| {
219 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
220 })
221 .map_err(sql_err)?;
222 let mut book = ContactBook::default();
223 for row in rows {
224 let (card_s, trust_s) = row.map_err(sql_err)?;
225 if let Ok(card) = ContactCard::from_bytes(card_s.as_bytes()) {
226 let trust = parse_trust(&trust_s);
227 let _ = book.insert(Contact { card, trust });
228 }
229 }
230 drop(stmt);
231 let mut astmt = db
232 .prepare("SELECT peer_id, alias FROM aliases")
233 .map_err(sql_err)?;
234 let arows = astmt
235 .query_map([], |row| {
236 Ok((row.get::<_, Vec<u8>>(0)?, row.get::<_, String>(1)?))
237 })
238 .map_err(sql_err)?;
239 for row in arows {
240 let (pid, alias) = row.map_err(sql_err)?;
241 if pid.len() == 32 {
242 let mut d = [0u8; 32];
243 d.copy_from_slice(&pid);
244 let id = PeerId::from_digest(d);
245 if let Some(c) = book.get_mut(&id) {
246 c.card.name = Some(alias);
247 }
248 }
249 }
250 Ok(book)
251 }
252
253 pub fn put_object(
254 &self,
255 env: &DropEnvelope,
256 manifest: &Manifest,
257 chunks: &[(u32, Vec<u8>)],
258 ownership: Ownership,
259 now: u64,
260 ) -> Result<ObjectId> {
261 verify_envelope(env, now)?;
262 let used = self.physical_bytes()?;
263 let incoming: u64 = chunks.iter().map(|(_, d)| d.len() as u64).sum();
264 if used.saturating_add(incoming) > self.quotas.maximum {
265 return Err(DdError::protocol(ErrorCode::Dds2001StoreFull, "quota"));
266 }
267 let env_bytes = encode_cbor(env)?;
268 let man_bytes = encode_cbor(manifest)?;
269 let oid = env.object_id;
270 {
271 let db = self.db.lock().expect("db");
272 db.execute(
273 "INSERT OR REPLACE INTO objects(object_id,envelope,source,created_at,expires_at,priority,hop_count,hop_limit,state,ownership,size,replication,application)
274 VALUES(?1,?2,?3,?4,?5,?6,?7,?8,?9,?10,?11,?12,?13)",
275 params![
276 oid.as_bytes().as_slice(),
277 env_bytes,
278 env.source.as_bytes().as_slice(),
279 env.creation_time as i64,
280 env.expiration as i64,
281 env.priority.as_u8() as i64,
282 env.hop_count as i64,
283 env.hop_limit as i64,
284 "stored",
285 ownership_str(ownership),
286 manifest.total_length as i64,
287 env.routing_policy.replication_budget as i64,
288 env.application,
289 ],
290 )
291 .map_err(sql_err)?;
292 db.execute(
293 "INSERT OR REPLACE INTO manifests(manifest_id, object_id, body) VALUES(?1,?2,?3)",
294 params![
295 manifest.id().as_bytes().as_slice(),
296 oid.as_bytes().as_slice(),
297 man_bytes
298 ],
299 )
300 .map_err(sql_err)?;
301 }
302 for (idx, refer) in manifest.chunks.iter().enumerate() {
303 self.db
304 .lock()
305 .expect("db")
306 .execute(
307 "INSERT OR REPLACE INTO object_chunks(object_id, idx, content_id, present) VALUES(?1,?2,?3,0)",
308 params![
309 oid.as_bytes().as_slice(),
310 idx as i64,
311 refer.id.as_bytes().as_slice()
312 ],
313 )
314 .map_err(sql_err)?;
315 }
316 for (idx, data) in chunks {
317 self.put_chunk(oid, *idx, data)?;
318 }
319 self.trace(oid, now, "stored")?;
320 Ok(oid)
321 }
322
323 pub fn put_chunk(&self, object_id: ObjectId, index: u32, data: &[u8]) -> Result<()> {
324 let refer = self.chunk_ref(object_id, index)?;
325 verify_chunk(&refer, data)?;
326 let path = self.chunk_path(&refer);
327 if !path.exists() {
328 fs::write(&path, data)?;
329 self.db
330 .lock()
331 .expect("db")
332 .execute(
333 "INSERT INTO chunks(content_id,size,refcount) VALUES(?1,?2,1)
334 ON CONFLICT(content_id) DO UPDATE SET refcount = refcount + 1",
335 params![refer.as_bytes().as_slice(), data.len() as i64],
336 )
337 .map_err(sql_err)?;
338 }
339 self.db
340 .lock()
341 .expect("db")
342 .execute(
343 "UPDATE object_chunks SET present=1 WHERE object_id=?1 AND idx=?2",
344 params![object_id.as_bytes().as_slice(), index as i64],
345 )
346 .map_err(sql_err)?;
347 Ok(())
348 }
349
350 fn chunk_ref(&self, object_id: ObjectId, index: u32) -> Result<crate::ChunkId> {
351 let db = self.db.lock().expect("db");
352 let bytes: Vec<u8> = db
353 .query_row(
354 "SELECT content_id FROM object_chunks WHERE object_id=?1 AND idx=?2",
355 params![object_id.as_bytes().as_slice(), index as i64],
356 |r| r.get(0),
357 )
358 .map_err(|_| DdError::protocol(ErrorCode::Dds2003MissingChunk, "index"))?;
359 if bytes.len() != 32 {
360 return Err(DdError::protocol(
361 ErrorCode::Dds2004CorruptMetadata,
362 "chunk id",
363 ));
364 }
365 let mut d = [0u8; 32];
366 d.copy_from_slice(&bytes);
367 Ok(crate::ChunkId::blake3(d))
368 }
369
370 fn chunk_path(&self, id: &crate::ChunkId) -> PathBuf {
371 self.root.join("chunks").join(hex_encode(id.as_bytes()))
372 }
373
374 pub fn get_envelope(&self, id: &ObjectId) -> Result<Option<DropEnvelope>> {
375 let db = self.db.lock().expect("db");
376 let blob: Option<Vec<u8>> = db
377 .query_row(
378 "SELECT envelope FROM objects WHERE object_id=?1",
379 params![id.as_bytes().as_slice()],
380 |r| r.get(0),
381 )
382 .optional()
383 .map_err(sql_err)?;
384 blob.map(|b| decode_cbor(&b)).transpose()
385 }
386
387 pub fn get_manifest(&self, id: &ObjectId) -> Result<Option<Manifest>> {
388 let db = self.db.lock().expect("db");
389 let blob: Option<Vec<u8>> = db
390 .query_row(
391 "SELECT body FROM manifests WHERE object_id=?1",
392 params![id.as_bytes().as_slice()],
393 |r| r.get(0),
394 )
395 .optional()
396 .map_err(sql_err)?;
397 blob.map(|b| decode_cbor(&b)).transpose()
398 }
399
400 pub fn present_mask(&self, id: &ObjectId) -> Result<Vec<bool>> {
401 let db = self.db.lock().expect("db");
402 let mut stmt = db
403 .prepare("SELECT idx, present FROM object_chunks WHERE object_id=?1 ORDER BY idx")
404 .map_err(sql_err)?;
405 let rows = stmt
406 .query_map(params![id.as_bytes().as_slice()], |r| {
407 Ok((r.get::<_, i64>(0)?, r.get::<_, i64>(1)?))
408 })
409 .map_err(sql_err)?;
410 let mut mask = Vec::new();
411 for row in rows {
412 let (idx, p) = row.map_err(sql_err)?;
413 let i = idx as usize;
414 if mask.len() <= i {
415 mask.resize(i + 1, false);
416 }
417 mask[i] = p != 0;
418 }
419 Ok(mask)
420 }
421
422 pub fn load_chunks(&self, id: &ObjectId) -> Result<Vec<Vec<u8>>> {
423 let man = self
424 .get_manifest(id)?
425 .ok_or_else(|| DdError::protocol(ErrorCode::Dds2003MissingChunk, "manifest"))?;
426 let slots = self.load_chunk_slots(id)?;
427 if man.erasure.is_some() {
428 return crate::chunk::reconstruct(&man, &slots);
429 }
430 let mut out = Vec::new();
431 for (i, refer) in man.chunks.iter().enumerate() {
432 let data = slots.get(i).and_then(|s| s.as_ref()).ok_or_else(|| {
433 DdError::protocol(ErrorCode::Dds2003MissingChunk, format!("chunk {i}"))
434 })?;
435 verify_chunk(&refer.id, data)?;
436 out.push(data.clone());
437 }
438 Ok(out)
439 }
440
441 pub fn load_chunk(&self, id: &ObjectId, index: u32) -> Result<Vec<u8>> {
442 let refer = self.chunk_ref(*id, index)?;
443 let data = fs::read(self.chunk_path(&refer))?;
444 verify_chunk(&refer, &data)?;
445 Ok(data)
446 }
447
448 pub fn inventory(&self, now: u64) -> Result<Vec<ObjectId>> {
449 let db = self.db.lock().expect("db");
450 let mut stmt = db
451 .prepare("SELECT object_id FROM objects WHERE expires_at > ?1 AND state != 'garbagecollected'")
452 .map_err(sql_err)?;
453 let rows = stmt
454 .query_map(params![now as i64], |r| r.get::<_, Vec<u8>>(0))
455 .map_err(sql_err)?;
456 let mut ids = Vec::new();
457 for row in rows {
458 let b = row.map_err(sql_err)?;
459 if b.len() == 32 {
460 let mut d = [0u8; 32];
461 d.copy_from_slice(&b);
462 ids.push(ObjectId::blake3(d));
463 }
464 }
465 ids.sort_by(|a, b| a.as_bytes().cmp(b.as_bytes()));
466 Ok(ids)
467 }
468
469 pub fn complete(&self, id: &ObjectId) -> Result<bool> {
470 let Some(man) = self.get_manifest(id)? else {
471 return Ok(false);
472 };
473 let mask = self.present_mask(id)?;
474 Ok(crate::chunk::can_recover(&man, &mask))
475 }
476
477 pub fn load_chunk_slots(&self, id: &ObjectId) -> Result<Vec<Option<Vec<u8>>>> {
478 let man = self
479 .get_manifest(id)?
480 .ok_or_else(|| DdError::protocol(ErrorCode::Dds2003MissingChunk, "manifest"))?;
481 let mut out = Vec::with_capacity(man.chunks.len());
482 for refer in &man.chunks {
483 let path = self.chunk_path(&refer.id);
484 if path.exists() {
485 let data = fs::read(&path)?;
486 if crate::chunk::verify_chunk(&refer.id, &data).is_ok() {
487 out.push(Some(data));
488 continue;
489 }
490 }
491 out.push(None);
492 }
493 Ok(out)
494 }
495
496 pub fn stats(&self) -> Result<StoreStats> {
497 let db = self.db.lock().expect("db");
498 let objects: i64 = db
499 .query_row("SELECT COUNT(*) FROM objects", [], |r| r.get(0))
500 .map_err(sql_err)?;
501 let chunks: i64 = db
502 .query_row("SELECT COUNT(*) FROM chunks", [], |r| r.get(0))
503 .map_err(sql_err)?;
504 let logical: i64 = db
505 .query_row("SELECT COALESCE(SUM(size),0) FROM objects", [], |r| {
506 r.get(0)
507 })
508 .map_err(sql_err)?;
509 let physical: i64 = db
510 .query_row("SELECT COALESCE(SUM(size),0) FROM chunks", [], |r| r.get(0))
511 .map_err(sql_err)?;
512 drop(db);
513 Ok(StoreStats {
514 objects: objects as u64,
515 payloads: objects as u64,
516 chunks: chunks as u64,
517 logical_size: logical as u64,
518 physical_size: physical as u64,
519 })
520 }
521
522 pub fn physical_bytes(&self) -> Result<u64> {
523 Ok(self.stats()?.physical_size)
524 }
525
526 pub fn gc(&self, now: u64) -> Result<u32> {
527 let db = self.db.lock().expect("db");
528 let n = db
529 .execute(
530 "UPDATE objects SET state='expired' WHERE expires_at <= ?1 AND ownership != 'local' AND object_id NOT IN (SELECT object_id FROM pins)",
531 params![now as i64],
532 )
533 .map_err(sql_err)?;
534 let n2 = db
535 .execute(
536 "UPDATE objects SET state='garbagecollected' WHERE state='expired'",
537 [],
538 )
539 .map_err(sql_err)?;
540 Ok((n + n2) as u32)
541 }
542
543 pub fn verify(&self) -> Result<Vec<String>> {
544 let mut problems = Vec::new();
545 let ids = self.inventory(unix_now()).unwrap_or_default();
546 for id in ids {
547 if let Ok(Some(man)) = self.get_manifest(&id) {
548 let mask = self.present_mask(&id)?;
549 for (i, refer) in man.chunks.iter().enumerate() {
550 if mask.get(i) == Some(&true) {
551 let path = self.chunk_path(&refer.id);
552 match fs::read(&path) {
553 Ok(d) => {
554 if verify_chunk(&refer.id, &d).is_err() {
555 problems.push(format!("corrupt chunk {i} of {id}"));
556 }
557 }
558 Err(_) => problems.push(format!("missing chunk file {i} of {id}")),
559 }
560 }
561 }
562 }
563 }
564 Ok(problems)
565 }
566
567 pub fn record_encounter(
568 &self,
569 peer: PeerId,
570 now: u64,
571 sent: u64,
572 recv: u64,
573 ok: bool,
574 ) -> Result<()> {
575 let db = self.db.lock().expect("db");
576 db.execute(
577 "INSERT INTO encounters(peer_id,first_seen,last_seen,encounter_count,bytes_sent,bytes_received,successful_forwards,failed_forwards)
578 VALUES(?1,?2,?2,1,?3,?4,?5,?6)
579 ON CONFLICT(peer_id) DO UPDATE SET
580 last_seen=excluded.last_seen,
581 encounter_count=encounter_count+1,
582 bytes_sent=bytes_sent+excluded.bytes_sent,
583 bytes_received=bytes_received+excluded.bytes_received,
584 successful_forwards=successful_forwards+excluded.successful_forwards,
585 failed_forwards=failed_forwards+excluded.failed_forwards",
586 params![
587 peer.as_bytes().as_slice(),
588 now as i64,
589 sent as i64,
590 recv as i64,
591 if ok { 1 } else { 0 },
592 if ok { 0 } else { 1 }
593 ],
594 )
595 .map_err(sql_err)?;
596 drop(db);
597 self.touch_hour_hist(peer, now)?;
598 Ok(())
599 }
600
601 pub fn encounter(&self, peer: PeerId) -> Result<Option<EncounterRow>> {
602 let db = self.db.lock().expect("db");
603 db.query_row(
604 "SELECT first_seen,last_seen,encounter_count,bytes_sent,bytes_received,successful_forwards,failed_forwards FROM encounters WHERE peer_id=?1",
605 params![peer.as_bytes().as_slice()],
606 |r| {
607 Ok(EncounterRow {
608 first_seen: r.get::<_, i64>(0)? as u64,
609 last_seen: r.get::<_, i64>(1)? as u64,
610 encounter_count: r.get::<_, i64>(2)? as u64,
611 bytes_sent: r.get::<_, i64>(3)? as u64,
612 bytes_received: r.get::<_, i64>(4)? as u64,
613 successful_forwards: r.get::<_, i64>(5)? as u64,
614 failed_forwards: r.get::<_, i64>(6)? as u64,
615 })
616 },
617 )
618 .optional()
619 .map_err(sql_err)
620 }
621
622 pub fn all_encounters(&self) -> Result<Vec<(PeerId, EncounterRow)>> {
623 let db = self.db.lock().expect("db");
624 let mut stmt = db
625 .prepare("SELECT peer_id,first_seen,last_seen,encounter_count,bytes_sent,bytes_received,successful_forwards,failed_forwards FROM encounters")
626 .map_err(sql_err)?;
627 let rows = stmt
628 .query_map([], |r| {
629 let b: Vec<u8> = r.get(0)?;
630 Ok((
631 b,
632 EncounterRow {
633 first_seen: r.get::<_, i64>(1)? as u64,
634 last_seen: r.get::<_, i64>(2)? as u64,
635 encounter_count: r.get::<_, i64>(3)? as u64,
636 bytes_sent: r.get::<_, i64>(4)? as u64,
637 bytes_received: r.get::<_, i64>(5)? as u64,
638 successful_forwards: r.get::<_, i64>(6)? as u64,
639 failed_forwards: r.get::<_, i64>(7)? as u64,
640 },
641 ))
642 })
643 .map_err(sql_err)?;
644 let mut out = Vec::new();
645 for row in rows {
646 let (b, e) = row.map_err(sql_err)?;
647 if b.len() == 32 {
648 let mut d = [0u8; 32];
649 d.copy_from_slice(&b);
650 out.push((PeerId::from_digest(d), e));
651 }
652 }
653 Ok(out)
654 }
655
656 pub fn set_state(&self, id: &ObjectId, state: DropState) -> Result<()> {
657 self.db
658 .lock()
659 .expect("db")
660 .execute(
661 "UPDATE objects SET state=?1 WHERE object_id=?2",
662 params![
663 format!("{state:?}").to_lowercase(),
664 id.as_bytes().as_slice()
665 ],
666 )
667 .map_err(sql_err)?;
668 Ok(())
669 }
670
671 pub fn decrement_replication(&self, id: &ObjectId) -> Result<u32> {
672 let db = self.db.lock().expect("db");
673 db.execute(
674 "UPDATE objects SET replication = MAX(replication-1,0) WHERE object_id=?1",
675 params![id.as_bytes().as_slice()],
676 )
677 .map_err(sql_err)?;
678 let v: i64 = db
679 .query_row(
680 "SELECT replication FROM objects WHERE object_id=?1",
681 params![id.as_bytes().as_slice()],
682 |r| r.get(0),
683 )
684 .map_err(sql_err)?;
685 Ok(v as u32)
686 }
687
688 pub fn replication(&self, id: &ObjectId) -> Result<u32> {
689 let v: i64 = self
690 .db
691 .lock()
692 .expect("db")
693 .query_row(
694 "SELECT replication FROM objects WHERE object_id=?1",
695 params![id.as_bytes().as_slice()],
696 |r| r.get(0),
697 )
698 .map_err(sql_err)?;
699 Ok(v as u32)
700 }
701
702 pub fn ownership(&self, id: &ObjectId) -> Result<Ownership> {
703 let s: String = self
704 .db
705 .lock()
706 .expect("db")
707 .query_row(
708 "SELECT ownership FROM objects WHERE object_id=?1",
709 params![id.as_bytes().as_slice()],
710 |r| r.get(0),
711 )
712 .map_err(sql_err)?;
713 Ok(parse_own(&s))
714 }
715
716 pub fn bump_hop(&self, env: &mut DropEnvelope) {
717 env.hop_count = env.hop_count.saturating_add(1);
718 }
719
720 pub fn trace(&self, id: ObjectId, ts: u64, event: &str) -> Result<()> {
721 self.db
722 .lock()
723 .expect("db")
724 .execute(
725 "INSERT INTO traces(object_id,ts,event) VALUES(?1,?2,?3)",
726 params![id.as_bytes().as_slice(), ts as i64, event],
727 )
728 .map_err(sql_err)?;
729 Ok(())
730 }
731
732 pub fn traces(&self, id: &ObjectId) -> Result<Vec<(u64, String)>> {
733 let db = self.db.lock().expect("db");
734 let mut stmt = db
735 .prepare("SELECT ts,event FROM traces WHERE object_id=?1 ORDER BY id")
736 .map_err(sql_err)?;
737 let rows = stmt
738 .query_map(params![id.as_bytes().as_slice()], |r| {
739 Ok((r.get::<_, i64>(0)? as u64, r.get::<_, String>(1)?))
740 })
741 .map_err(sql_err)?;
742 let mut v = Vec::new();
743 for row in rows {
744 v.push(row.map_err(sql_err)?);
745 }
746 Ok(v)
747 }
748
749 pub fn pin(&self, id: &ObjectId) -> Result<()> {
750 self.db
751 .lock()
752 .expect("db")
753 .execute(
754 "INSERT OR IGNORE INTO pins(object_id) VALUES(?1)",
755 params![id.as_bytes().as_slice()],
756 )
757 .map_err(sql_err)?;
758 Ok(())
759 }
760
761 pub fn unpin(&self, id: &ObjectId) -> Result<()> {
762 self.db
763 .lock()
764 .expect("db")
765 .execute(
766 "DELETE FROM pins WHERE object_id=?1",
767 params![id.as_bytes().as_slice()],
768 )
769 .map_err(sql_err)?;
770 Ok(())
771 }
772
773 pub fn is_pinned(&self, id: &ObjectId) -> Result<bool> {
774 let n: i64 = self
775 .db
776 .lock()
777 .expect("db")
778 .query_row(
779 "SELECT COUNT(*) FROM pins WHERE object_id=?1",
780 params![id.as_bytes().as_slice()],
781 |r| r.get(0),
782 )
783 .map_err(sql_err)?;
784 Ok(n > 0)
785 }
786
787 pub fn tag(&self, id: &ObjectId, tag: &str) -> Result<()> {
788 self.db
789 .lock()
790 .expect("db")
791 .execute(
792 "INSERT OR IGNORE INTO tags(object_id, tag) VALUES(?1,?2)",
793 params![id.as_bytes().as_slice(), tag],
794 )
795 .map_err(sql_err)?;
796 Ok(())
797 }
798
799 pub fn set_alias(&self, peer: PeerId, alias: &str) -> Result<()> {
800 self.db
801 .lock()
802 .expect("db")
803 .execute(
804 "INSERT OR REPLACE INTO aliases(peer_id, alias) VALUES(?1,?2)",
805 params![peer.as_bytes().as_slice(), alias],
806 )
807 .map_err(sql_err)?;
808 Ok(())
809 }
810
811 pub fn alias(&self, peer: PeerId) -> Result<Option<String>> {
812 self.db
813 .lock()
814 .expect("db")
815 .query_row(
816 "SELECT alias FROM aliases WHERE peer_id=?1",
817 params![peer.as_bytes().as_slice()],
818 |r| r.get(0),
819 )
820 .optional()
821 .map_err(sql_err)
822 }
823
824 pub fn search_meta(&self, q: &str) -> Result<Vec<ObjectId>> {
825 let needle = q.to_ascii_lowercase();
826 let mut hits = Vec::new();
827 for id in self.inventory(unix_now())? {
828 let Some(env) = self.get_envelope(&id)? else {
829 continue;
830 };
831 let blob = format!(
832 "{} {} {:?} {}",
833 env.application,
834 env.source,
835 env.destination,
836 env.topic.as_deref().unwrap_or("")
837 )
838 .to_ascii_lowercase();
839 if blob.contains(&needle) {
840 hits.push(id);
841 }
842 }
843 Ok(hits)
844 }
845
846 pub fn put_space(&self, rec: &crate::space::SpaceRecord) -> Result<()> {
847 let body = serde_json::to_string(rec).map_err(|e| DdError::crypto(e.to_string()))?;
848 self.db
849 .lock()
850 .expect("db")
851 .execute(
852 "INSERT OR REPLACE INTO spaces(name, body) VALUES(?1,?2)",
853 params![rec.name, body],
854 )
855 .map_err(sql_err)?;
856 Ok(())
857 }
858
859 pub fn list_spaces(&self) -> Result<Vec<crate::space::SpaceRecord>> {
860 let db = self.db.lock().expect("db");
861 let mut stmt = db.prepare("SELECT body FROM spaces").map_err(sql_err)?;
862 let rows = stmt
863 .query_map([], |r| r.get::<_, String>(0))
864 .map_err(sql_err)?;
865 let mut out = Vec::new();
866 for row in rows {
867 if let Ok(s) = serde_json::from_str(&row.map_err(sql_err)?) {
868 out.push(s);
869 }
870 }
871 Ok(out)
872 }
873
874 pub fn subscribe_channel(&self, sub: &crate::channel::ChannelSub) -> Result<()> {
875 self.db
876 .lock()
877 .expect("db")
878 .execute(
879 "INSERT OR REPLACE INTO channels(name, filter, policy) VALUES(?1,?2,?3)",
880 params![
881 sub.name,
882 format!("{:?}", sub.filter).to_lowercase(),
883 format!("{:?}", sub.policy).to_lowercase()
884 ],
885 )
886 .map_err(sql_err)?;
887 Ok(())
888 }
889
890 pub fn list_channels(&self) -> Result<Vec<(String, String, String)>> {
891 let db = self.db.lock().expect("db");
892 let mut stmt = db
893 .prepare("SELECT name, filter, policy FROM channels")
894 .map_err(sql_err)?;
895 let rows = stmt
896 .query_map([], |r| {
897 Ok((
898 r.get::<_, String>(0)?,
899 r.get::<_, String>(1)?,
900 r.get::<_, String>(2)?,
901 ))
902 })
903 .map_err(sql_err)?;
904 let mut out = Vec::new();
905 for row in rows {
906 out.push(row.map_err(sql_err)?);
907 }
908 Ok(out)
909 }
910
911 pub fn put_circle(&self, name: &str, members: &[PeerId]) -> Result<()> {
912 let ids: Vec<String> = members.iter().map(|p| p.to_string()).collect();
913 let body = ids.join(",");
914 self.db
915 .lock()
916 .expect("db")
917 .execute(
918 "INSERT OR REPLACE INTO circles(name, members) VALUES(?1,?2)",
919 params![name, body],
920 )
921 .map_err(sql_err)?;
922 Ok(())
923 }
924
925 pub fn circle_members(&self, name: &str) -> Result<Vec<PeerId>> {
926 let s: String = self
927 .db
928 .lock()
929 .expect("db")
930 .query_row(
931 "SELECT members FROM circles WHERE name=?1",
932 params![name],
933 |r| r.get(0),
934 )
935 .map_err(sql_err)?;
936 let mut out = Vec::new();
937 for p in s.split(',') {
938 if let Ok(id) = p.parse() {
939 out.push(id);
940 }
941 }
942 Ok(out)
943 }
944
945 pub fn compact(&self) -> Result<u32> {
946 let n = self.gc(unix_now())?;
947 let _ = self
948 .db
949 .lock()
950 .expect("db")
951 .execute_batch("PRAGMA wal_checkpoint(TRUNCATE); VACUUM;");
952 Ok(n)
953 }
954
955 pub fn quotas(&self) -> &StorageQuotas {
956 &self.quotas
957 }
958
959 pub fn kv_get(&self, k: &str) -> Result<Option<String>> {
960 let db = self.db.lock().expect("db");
961 db.query_row("SELECT v FROM kv WHERE k=?1", params![k], |r| r.get(0))
962 .optional()
963 .map_err(sql_err)
964 }
965
966 pub fn kv_set(&self, k: &str, v: &str) -> Result<()> {
967 self.db
968 .lock()
969 .expect("db")
970 .execute(
971 "INSERT OR REPLACE INTO kv(k,v) VALUES(?1,?2)",
972 params![k, v],
973 )
974 .map_err(sql_err)?;
975 Ok(())
976 }
977
978 fn touch_hour_hist(&self, peer: PeerId, now: u64) -> Result<()> {
979 let key = format!("eh:{}", crate::hex_encode(peer.as_bytes()));
980 let mut hist = [0u32; 24];
981 if let Some(s) = self.kv_get(&key)?
982 && let Ok(v) = serde_json::from_str::<Vec<u32>>(&s)
983 {
984 for (i, n) in v.into_iter().take(24).enumerate() {
985 hist[i] = n;
986 }
987 }
988 let hour = ((now % 86_400) / 3600) as usize;
989 hist[hour] = hist[hour].saturating_add(1);
990 self.kv_set(
991 &key,
992 &serde_json::to_string(&hist.to_vec()).unwrap_or_else(|_| "[]".into()),
993 )
994 }
995
996 pub fn hour_hist(&self, peer: PeerId) -> Result<[u32; 24]> {
997 let key = format!("eh:{}", crate::hex_encode(peer.as_bytes()));
998 let mut hist = [0u32; 24];
999 if let Some(s) = self.kv_get(&key)?
1000 && let Ok(v) = serde_json::from_str::<Vec<u32>>(&s)
1001 {
1002 for (i, n) in v.into_iter().take(24).enumerate() {
1003 hist[i] = n;
1004 }
1005 }
1006 Ok(hist)
1007 }
1008
1009 pub fn put_receipt(&self, r: &crate::receipt::Receipt) -> Result<()> {
1010 let body = crate::protocol::encode_cbor(r)?;
1011 let rid = DefaultProvider.hash(crate::HashAlgorithm::Blake3, &body);
1012 self.db
1013 .lock()
1014 .expect("db")
1015 .execute(
1016 "INSERT OR REPLACE INTO receipts(receipt_id, object_id, kind, body) VALUES(?1,?2,?3,?4)",
1017 params![
1018 rid.0.as_slice(),
1019 r.object_id.as_bytes().as_slice(),
1020 format!("{:?}", r.kind).to_lowercase(),
1021 body
1022 ],
1023 )
1024 .map_err(sql_err)?;
1025 Ok(())
1026 }
1027
1028 pub fn receipts_for(&self, id: &ObjectId) -> Result<Vec<crate::receipt::Receipt>> {
1029 let db = self.db.lock().expect("db");
1030 let mut stmt = db
1031 .prepare("SELECT body FROM receipts WHERE object_id=?1")
1032 .map_err(sql_err)?;
1033 let rows = stmt
1034 .query_map(params![id.as_bytes().as_slice()], |r| {
1035 r.get::<_, Vec<u8>>(0)
1036 })
1037 .map_err(sql_err)?;
1038 let mut out = Vec::new();
1039 for row in rows {
1040 let b = row.map_err(sql_err)?;
1041 if let Ok(r) = crate::protocol::decode_cbor(&b) {
1042 out.push(r);
1043 }
1044 }
1045 Ok(out)
1046 }
1047
1048 pub fn receipt_issuer_count(&self, id: &ObjectId) -> Result<u32> {
1049 let rs = self.receipts_for(id)?;
1050 let mut seen = std::collections::BTreeSet::new();
1051 for r in rs {
1052 seen.insert(*r.issuer.as_bytes());
1053 }
1054 Ok(seen.len() as u32)
1055 }
1056
1057 pub fn repair_safe(&self) -> Result<Vec<String>> {
1059 let mut notes = Vec::new();
1060 let n = self.compact()?;
1061 notes.push(format!("vacuum/gc rows={n} (local/pinned Drops kept)"));
1062 let chunk_dir = self.root.join("chunks");
1063 if let Ok(rd) = fs::read_dir(&chunk_dir) {
1064 let mut orphans = 0u32;
1065 for ent in rd.flatten() {
1066 let name = ent.file_name();
1067 let hex = name.to_string_lossy();
1068 let db = self.db.lock().expect("db");
1069 let n: i64 = db
1070 .query_row(
1071 "SELECT COUNT(*) FROM chunks WHERE lower(hex(content_id))=?1",
1072 params![hex.to_ascii_lowercase()],
1073 |r| r.get(0),
1074 )
1075 .unwrap_or(1);
1076 drop(db);
1077 if n == 0 {
1078 let tmp = ent.path();
1079 if tmp.extension().and_then(|e| e.to_str()) == Some("tmp") {
1080 let _ = fs::remove_file(&tmp);
1081 orphans += 1;
1082 }
1083 }
1084 }
1085 if orphans > 0 {
1086 notes.push(format!("removed {orphans} temporary chunk files"));
1087 }
1088 }
1089 for lock in ["daemon.pid.lock", ".write-probe", "control.sock.lock"] {
1090 let p = self.root.join(lock);
1091 if p.exists() {
1092 let _ = fs::remove_file(&p);
1093 notes.push(format!("removed stale {lock}"));
1094 }
1095 }
1096 Ok(notes)
1097 }
1098}
1099
1100#[derive(Debug, Clone)]
1101pub struct StoreStats {
1102 pub objects: u64,
1103 pub payloads: u64,
1104 pub chunks: u64,
1105 pub logical_size: u64,
1106 pub physical_size: u64,
1107}
1108
1109impl StoreStats {
1110 pub fn dedup_ratio(&self) -> f64 {
1111 if self.logical_size == 0 {
1112 0.0
1113 } else {
1114 1.0 - (self.physical_size as f64 / self.logical_size as f64)
1115 }
1116 }
1117}
1118
1119#[derive(Debug, Clone)]
1120pub struct EncounterRow {
1121 pub first_seen: u64,
1122 pub last_seen: u64,
1123 pub encounter_count: u64,
1124 pub bytes_sent: u64,
1125 pub bytes_received: u64,
1126 pub successful_forwards: u64,
1127 pub failed_forwards: u64,
1128}
1129
1130fn sql_err(e: rusqlite::Error) -> DdError {
1131 DdError::protocol(ErrorCode::Dds2004CorruptMetadata, e.to_string())
1132}
1133
1134fn ownership_str(o: Ownership) -> &'static str {
1135 match o {
1136 Ownership::Local => "local",
1137 Ownership::Incoming => "incoming",
1138 Ownership::Relay => "relay",
1139 Ownership::Temporary => "temporary",
1140 }
1141}
1142
1143fn parse_own(s: &str) -> Ownership {
1144 match s {
1145 "local" => Ownership::Local,
1146 "incoming" => Ownership::Incoming,
1147 "temporary" => Ownership::Temporary,
1148 _ => Ownership::Relay,
1149 }
1150}
1151
1152fn parse_trust(s: &str) -> TrustState {
1153 match s {
1154 "verified" => TrustState::Verified,
1155 "known" => TrustState::Known,
1156 "blocked" => TrustState::Blocked,
1157 "observed" => TrustState::Observed,
1158 _ => TrustState::Unknown,
1159 }
1160}
1161
1162pub fn unix_now() -> u64 {
1163 SystemTime::now()
1164 .duration_since(UNIX_EPOCH)
1165 .map(|d| d.as_secs())
1166 .unwrap_or(0)
1167}
1168
1169pub fn format_bytes(n: u64) -> String {
1170 const KB: f64 = 1024.0;
1171 let x = n as f64;
1172 if x >= KB * KB * KB {
1173 format!("{:.1} GB", x / (KB * KB * KB))
1174 } else if x >= KB * KB {
1175 format!("{:.1} MB", x / (KB * KB))
1176 } else if x >= KB {
1177 format!("{:.1} KB", x / KB)
1178 } else {
1179 format!("{n} B")
1180 }
1181}