Skip to main content

deaddrop_core/store/
mod.rs

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    /// Safe repair: indexes, WAL, unreferenced chunk files. Never deletes Drop objects.
1058    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}