diff --git a/src/storage.rs b/src/storage.rs index 8fae82a..a9867e3 100644 --- a/src/storage.rs +++ b/src/storage.rs @@ -340,8 +340,13 @@ impl Database { Ok(db) } + /// Erwirbt den Mutex auf die SQLite-Verbindung und fängt eventuelles Mutex-Poisoning ab (Z-02). + pub fn conn(&self) -> std::sync::MutexGuard<'_, Connection> { + self.conn.lock().unwrap_or_else(|e| e.into_inner()) + } + pub fn conn_for_test(&self) -> std::sync::MutexGuard<'_, Connection> { - self.conn.lock().unwrap() + self.conn() } /// Erzeugt eine geklonte Instanz mit einer isolierten aktiven Session (Slot & DEK). @@ -359,14 +364,14 @@ impl Database { /// Setzt den aktiven Slot und DEK für automatische Metadaten-Authentifizierung (K-01). pub fn set_active_slot_and_dek(&self, slot_id: u32, dek: Zeroizing<[u8; 32]>) { - *self.active_session.lock().unwrap() = Some((slot_id, dek)); + *self.active_session.lock().unwrap_or_else(|e| e.into_inner()) = Some((slot_id, dek)); } /// Gibt den aktuellen aktiven DEK zurück, falls gesetzt. pub fn active_dek(&self) -> Option> { self.active_session .lock() - .unwrap() + .unwrap_or_else(|e| e.into_inner()) .as_ref() .map(|(_, dek)| dek.clone()) } @@ -386,7 +391,7 @@ impl Database { /// Führt automatische, rückwärtskompatible Schema-Upgrades (z. B. Spalte slot_id, auto_vacuum, is_carrier) durch. pub fn ensure_schema_upgrades(&self) -> Result<()> { - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); let av: i64 = conn .query_row("PRAGMA auto_vacuum;", [], |r| r.get(0)) .unwrap_or(0); @@ -432,7 +437,7 @@ impl Database { /// Setzt die vorgeschriebenen SQLite3-Pragmas: 8192 Page-Size, Incremental Auto-Vacuum, WAL, NORMAL synchronous, Secure Delete. pub fn init_pragmas(&self) -> Result<()> { - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); conn.execute_batch( "PRAGMA page_size = 8192; PRAGMA auto_vacuum = INCREMENTAL; @@ -447,7 +452,7 @@ impl Database { /// Führt ein inkrementelles Auto-Vacuum aus, um freigegebene Datenbankseiten an das Betriebssystem zurückzugeben. pub fn incremental_vacuum(&self, pages: Option) -> Result { - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); let before: i64 = conn .query_row("PRAGMA freelist_count;", [], |r| r.get(0)) .unwrap_or(0); @@ -478,14 +483,14 @@ impl Database { /// Gibt die Anzahl ungenutzter Freelist-Seiten zurück. pub fn freelist_count(&self) -> Result { - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); let count: i64 = conn.query_row("PRAGMA freelist_count;", [], |r| r.get(0))?; Ok(count as usize) } /// Sucht nach einem existierenden Carrier-Knoten (is_carrier = 1) (S-03). pub fn find_carrier_node_id(&self) -> Result> { - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); let mut stmt = conn.prepare("SELECT id FROM nodes WHERE is_carrier = 1 LIMIT 1")?; let id = stmt.query_row([], |r| r.get::<_, i64>(0)).optional()?; Ok(id) @@ -493,7 +498,7 @@ impl Database { /// Markiert einen existierenden Knoten explizit als Carrier. pub fn mark_carrier_node_id(&self, id: i64) -> Result<()> { - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); conn.execute("UPDATE nodes SET is_carrier = 1 WHERE id = ?1", [id])?; Ok(()) } @@ -503,7 +508,7 @@ impl Database { if node_id == ancestor_id { return Ok(true); } - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); let mut stmt = conn.prepare( "WITH RECURSIVE sub(id) AS ( SELECT id FROM nodes WHERE id = ?1 @@ -529,7 +534,7 @@ impl Database { /// garantiert werden, da der Flash Translation Layer (FTL) und Wear-Leveling-Mechanismen Sektoren /// neuen Flash-Speicherzellen zuweisen. pub fn shred_chunks_for_node(&self, node_id: i64) -> Result<()> { - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); let mut stmt = conn.prepare("SELECT chunk_index, length(ciphertext) FROM chunks WHERE node_id = ?1")?; let chunks: Vec<(u32, usize)> = stmt @@ -579,7 +584,7 @@ impl Database { &[u8; 32], // raw DEK_1 )>, ) -> Result> { - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); conn.execute_batch( "CREATE TABLE IF NOT EXISTS meta ( @@ -789,7 +794,7 @@ impl Database { header_tag: &[u8; 16], hidden: Option<(&[u8; 16], &KdfParams, &[u8], &[u8; 12], &[u8; 16])>, ) -> Result<()> { - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); conn.execute_batch( "CREATE TABLE IF NOT EXISTS meta ( @@ -925,7 +930,7 @@ impl Database { /// Liest alle Header-Slots aus der `meta`-Tabelle aus (Slot 0 = Decoy/Standard, Slot 1 = Hidden Vault oder Dummy-Rauschen). /// SA-01: Führt strikte Vorab-Validierung der Container-Struktur VOR jeglicher KDF-Berechnung durch. pub fn read_slots(&self) -> Result> { - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); // SA-01 (HIGH): Vorab-Prüfung der Gesamtzahl der Slots (Schutz gegen KDF-Amplification / Container DoS) let slot_count: i64 = conn.query_row("SELECT count(*) FROM meta", [], |r| r.get(0))?; @@ -1046,7 +1051,7 @@ impl Database { } let slot0 = &slots[0]; - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); let magic: Vec = conn.query_row( "SELECT magic FROM meta WHERE slot_id = 0 LIMIT 1", [], @@ -1107,7 +1112,7 @@ impl Database { new_header_nonce: &[u8; 12], new_header_tag: &[u8; 16], ) -> Result<()> { - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); let params_json = serde_json::to_string(new_params)?; let rows_affected = conn.execute( "UPDATE meta SET kdf_salt = ?1, kdf_params = ?2, wrapped_dek = ?3, header_nonce = ?4, header_tag = ?5 WHERE slot_id = ?6", @@ -1133,7 +1138,7 @@ impl Database { /// Aktualisiert die Version in der meta-Tabelle (z. B. für Migrationen oder Tests). pub fn set_meta_version(&self, version: u32) -> Result<()> { - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); let rows_affected = conn.execute("UPDATE meta SET version = ?1", params![version])?; if rows_affected == 0 { bail!("Konnte Container-Version nicht aktualisieren: meta-Tabelle ist leer"); @@ -1158,7 +1163,7 @@ impl Database { new_dek: &[u8; 32], version: u32, ) -> Result<()> { - let mut conn = self.conn.lock().unwrap(); + let mut conn = self.conn(); let tx = conn.transaction()?; // Ermittle alle Node-IDs dieses Vaults @@ -1284,7 +1289,7 @@ impl Database { } let segments: Vec<&str> = normalized.split('/').filter(|s| !s.is_empty()).collect(); - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); let mut current_id = root_id; let mut last_record = None; @@ -1387,7 +1392,7 @@ impl Database { vault_id: u32, dek: &[u8; 32], ) -> Result> { - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); let mut stmt = conn.prepare( "SELECT id, parent_id, name, is_dir, size, created_at, modified_at FROM nodes WHERE id = ?1", @@ -1429,7 +1434,7 @@ impl Database { vault_id: u32, dek: &[u8; 32], ) -> Result> { - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); let mut stmt = conn.prepare( "SELECT id, parent_id, name, is_dir, size, created_at, modified_at FROM nodes @@ -1479,7 +1484,7 @@ impl Database { ) -> Result { crate::pathutil::validate_node_name(name)?; let now = current_timestamp(); - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); let stored_name = if vault_id == 1 { encrypt_node_name(dek, parent_id, name) @@ -1521,7 +1526,7 @@ impl Database { /// Aktualisiert Dateigröße und Modifikationszeitstempel eines Knotens. pub fn update_node_size_and_time(&self, id: i64, size: u64, modified_at: u64) -> Result<()> { - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); conn.execute( "UPDATE nodes SET size = ?1, modified_at = ?2 WHERE id = ?3", params![size as i64, modified_at as i64, id], @@ -1562,7 +1567,7 @@ impl Database { // 2. Shredde auch rekursiv alle Unterknoten let child_ids: Vec = { - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); let mut stmt = conn.prepare("SELECT id FROM nodes WHERE parent_id = ?1")?; let ids = stmt .query_map(params![id], |row| row.get(0))? @@ -1575,7 +1580,7 @@ impl Database { let _ = self.delete_node(child_id); } - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); conn.execute("DELETE FROM chunks WHERE node_id = ?1", params![id])?; conn.execute("DELETE FROM nodes WHERE id = ?1", params![id])?; drop(conn); @@ -1595,7 +1600,7 @@ impl Database { self.assert_not_carrier(id)?; crate::pathutil::validate_node_name(new_name)?; let now = current_timestamp(); - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); let stored_name = if vault_id == 1 { encrypt_node_name(dek, new_parent_id, new_name) } else { @@ -1617,7 +1622,7 @@ impl Database { /// Liest einen verschlüsselten Chunk aus der Datenbank. pub fn read_chunk(&self, node_id: i64, chunk_index: u32) -> Result> { - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); let mut stmt = conn.prepare( "SELECT nonce, tag, ciphertext, generation FROM chunks WHERE node_id = ?1 AND chunk_index = ?2", )?; @@ -1655,7 +1660,7 @@ impl Database { /// Ermittelt die nächste Generation für einen Chunk (K-02 Chunk-Replay-Schutz). /// Garantiert eine strikt monoton steigende Generation containerweit. pub fn next_chunk_generation(&self, node_id: i64, chunk_index: u32) -> Result { - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); let current_gen: Option = conn .query_row( "SELECT generation FROM chunks WHERE node_id = ?1 AND chunk_index = ?2", @@ -1685,7 +1690,7 @@ impl Database { tag: &[u8; 16], ciphertext: &[u8], ) -> Result<()> { - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); conn.execute( "INSERT INTO chunks (node_id, chunk_index, generation, nonce, tag, ciphertext) VALUES (?1, ?2, ?3, ?4, ?5, ?6) @@ -1720,7 +1725,7 @@ impl Database { modified_at: u64, ) -> Result<()> { self.assert_not_carrier(node_id)?; - let mut conn = self.conn.lock().unwrap(); + let mut conn = self.conn(); let tx = conn.transaction()?; tx.execute( "INSERT INTO chunks (node_id, chunk_index, generation, nonce, tag, ciphertext) @@ -1753,7 +1758,7 @@ impl Database { /// und shreddert die abzuschneidenden Chunks vorher mit kryptografischem Zufallsrauschen. pub fn truncate_chunks_after(&self, node_id: i64, max_chunk_index: u32) -> Result<()> { self.assert_not_carrier(node_id)?; - let mut conn = self.conn.lock().unwrap(); + let mut conn = self.conn(); let tx = conn.transaction()?; { let mut stmt = tx.prepare( @@ -1802,7 +1807,7 @@ impl Database { /// Löscht und shreddert alle Chunks eines Knotens (z. B. beim Kürzen auf 0 Bytes) (S-05). pub fn delete_all_chunks(&self, node_id: i64) -> Result<()> { self.assert_not_carrier(node_id)?; - let mut conn = self.conn.lock().unwrap(); + let mut conn = self.conn(); let tx = conn.transaction()?; { let mut stmt = tx @@ -1846,7 +1851,7 @@ impl Database { /// Prüft den aktuellen Advisory-Lock-Status (S-06). /// Gibt `Some((pid, host, timestamp))` zurück, falls ein Lock aktiv und der Prozess noch am Leben ist. pub fn check_advisory_lock(&self) -> Result> { - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); let mut stmt = conn .prepare("SELECT lock_pid, lock_host, lock_time FROM meta WHERE slot_id = 0 LIMIT 1")?; let lock_info = stmt @@ -1899,7 +1904,7 @@ impl Database { .unwrap_or_else(|_| "localhost".to_string()); let now = current_timestamp(); - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); conn.execute( "UPDATE meta SET lock_pid = ?1, lock_host = ?2, lock_time = ?3 WHERE slot_id = 0", params![pid as i64, host, now], @@ -1909,7 +1914,7 @@ impl Database { /// Entfernt den Advisory Lock (S-06). pub fn release_advisory_lock(&self) -> Result<()> { - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); conn.execute( "UPDATE meta SET lock_pid = NULL, lock_host = NULL, lock_time = NULL WHERE slot_id = 0", [], @@ -1930,7 +1935,7 @@ impl Database { /// Erzeugt die deterministische kanonische Byterepräsentation für einen spezifischen Vault. pub fn canonical_nodes_bytes_for_vault(&self, vault_id: u32) -> Result> { - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); let mut buf = Vec::new(); if vault_id == 1 { @@ -2009,12 +2014,12 @@ impl Database { /// Aktualisiert den Metadaten-MAC des aktiven Slots bei strukturellen Modifikationen (Format V3 / K-01). pub fn update_metadata_mac(&self) -> Result<()> { - let session_opt = self.active_session.lock().unwrap().clone(); + let session_opt = self.active_session.lock().unwrap_or_else(|e| e.into_inner()).clone(); let Some((slot_id, dek)) = session_opt else { return Ok(()); }; - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); let version_and_gen: Option<(u32, u64)> = conn .query_row( "SELECT version, metadata_gen FROM meta WHERE slot_id = ?1 LIMIT 1", @@ -2038,7 +2043,7 @@ impl Database { let mac_key = derive_metadata_mac_key(&dek); let new_mac = compute_metadata_mac(&mac_key, next_gen, &canonical); - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); conn.execute( "UPDATE meta SET metadata_mac = ?1, metadata_gen = ?2 WHERE slot_id = ?3", params![new_mac.as_slice(), next_gen, slot_id], @@ -2060,7 +2065,7 @@ impl Database { /// Prüft die Integrität des Metadaten-MAC für einen spezifischen Slot (0: Decoy, 1: Hidden). pub fn verify_metadata_mac_for_slot(&self, slot_id: u32, dek: &[u8; 32]) -> Result { - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); let meta_row: Option<(u32, Option>, u64)> = conn .query_row( "SELECT version, metadata_mac, metadata_gen FROM meta WHERE slot_id = ?1 LIMIT 1", @@ -2101,7 +2106,7 @@ impl Database { /// Führt ein Upgrade des Containerformats auf Format V3 durch (Format V3 / K-01 & K-02). pub fn upgrade_to_v3(&self, dek: &[u8; 32]) -> Result<()> { - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); let version: u32 = conn.query_row( "SELECT version FROM meta WHERE slot_id = 0 LIMIT 1", [], @@ -2195,7 +2200,7 @@ impl Database { /// Erzwingt einen SQLite WAL Checkpoint und leert das Write-Ahead-Log. pub fn checkpoint(&self) -> Result<()> { - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); let _res: (i64, i64, i64) = conn.query_row("PRAGMA wal_checkpoint(TRUNCATE);", [], |row| { Ok((row.get(0)?, row.get(1)?, row.get(2)?)) @@ -2215,7 +2220,7 @@ impl Database { } let mut dest_conn = Connection::open(dest_path)?; - let src_conn = self.conn.lock().unwrap(); + let src_conn = self.conn(); let backup = rusqlite::backup::Backup::new(&src_conn, &mut dest_conn)?; backup.run_to_completion(100, std::time::Duration::from_millis(20), None)?; @@ -2261,7 +2266,7 @@ impl Database { /// Schreibt oder stellt die Metadaten in der `meta`-Tabelle wieder her (z. B. nach Restore oder Header-Neugenerierung). pub fn restore_meta(&self, meta: &ContainerMeta) -> Result<()> { - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); conn.execute_batch( "CREATE TABLE IF NOT EXISTS meta ( @@ -2341,7 +2346,7 @@ impl Database { /// Führt SQLite-eigene Integritäts- und Foreign-Key-Prüfungen aus. pub fn run_sqlite_integrity_check(&self) -> Result> { - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); let mut issues = Vec::new(); // 1. PRAGMA integrity_check @@ -2374,7 +2379,7 @@ impl Database { /// Zählt die Anzahl von Verzeichnissen, Dateien und Daten-Chunks im Container. pub fn count_nodes_and_chunks(&self) -> Result<(usize, usize, usize)> { - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); let dirs: i64 = conn.query_row("SELECT COUNT(*) FROM nodes WHERE is_dir = 1", [], |r| { r.get(0) })?; @@ -2388,7 +2393,7 @@ impl Database { /// Listet alle Knoten (Dateien und Ordner) im gesamten Baum auf. pub fn list_all_nodes(&self) -> Result> { - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); let mut stmt = conn.prepare( "SELECT id, parent_id, name, is_dir, size, created_at, modified_at FROM nodes ORDER BY id ASC", )?; @@ -2412,7 +2417,7 @@ impl Database { /// Liefert alle vorhandenen Chunk-Identifikatoren (node_id, chunk_index). pub fn list_all_chunk_headers(&self) -> Result> { - let conn = self.conn.lock().unwrap(); + let conn = self.conn(); let mut stmt = conn.prepare("SELECT node_id, chunk_index FROM chunks ORDER BY node_id, chunk_index")?; let rows = stmt.query_map([], |row| { @@ -2720,7 +2725,7 @@ mod tests { ); // Forensische Prüfung: Roh-Inspektion der SQLite-Tabellen - let conn = db.conn.lock().unwrap(); + let conn = db.conn(); let raw_name_v0: String = conn .query_row( "SELECT name FROM nodes WHERE id = ?1", @@ -2822,4 +2827,56 @@ mod tests { assert_eq!(updated_file.size, new_size); assert_eq!(updated_file.modified_at, modified_at); } + + #[test] + fn test_z02_poisoned_sqlite_mutex_recovery() { + let db = Database::open_in_memory().unwrap(); + let (salt, kdf, wrapped_dek, nonce, tag) = ( + [1u8; 16], + KdfParams::default(), + vec![2u8; 40], + [3u8; 12], + [4u8; 16], + ); + db.init_schema(&salt, &kdf, &wrapped_dek, &nonce, &tag) + .unwrap(); + + // 1. Verifiziere normale Funktion + let node1 = db.create_node(1, "normal.txt", false).unwrap(); + assert_eq!(node1.name, "normal.txt"); + + // 2. Simuliere Thread-Panic während gehaltener Mutex-Sperre auf self.conn + let db_clone = db.clone(); + let handle = std::thread::spawn(move || { + let _guard = db_clone.conn(); + panic!("Simulierter Crash im Worker-Thread während aktiver SQLite-Verbindung"); + }); + let res = handle.join(); + assert!(res.is_err(), "Worker-Thread muss wie erwartet gepanict haben"); + + // 3. Mutex ist nun poisoned. Ohne Z-02 schlägt jeder nachfolgende Aufruf fehl. + // Mit Z-02 fängt unwrap_or_else(|e| e.into_inner()) das Poisoning ab: + let node2 = db.create_node(1, "after_poison.txt", false); + assert!( + node2.is_ok(), + "Z-02: Datenbankoperationen müssen trotz vergiftetem Mutex erfolgreich wiederhergestellt werden" + ); + assert_eq!(node2.unwrap().name, "after_poison.txt"); + + // 4. Prüfe auch active_session Poisoning + let db_clone2 = db.clone(); + let handle2 = std::thread::spawn(move || { + let _guard = db_clone2.active_session.lock().unwrap(); + panic!("Simulierter Crash während active_session Sperre"); + }); + let _ = handle2.join(); + + // Setzen und Lesen muss weiterhin fehlerfrei funktionieren + let test_dek = Zeroizing::new([42u8; 32]); + db.set_active_dek(test_dek); + let read_dek = db.active_dek(); + assert!(read_dek.is_some()); + assert_eq!(read_dek.unwrap()[0], 42); + } } +