Skip to main content

volli_core/
peer_db.rs

1//! SQLite-based peer database management
2//!
3//! This module provides persistent storage for peer information using SQLite
4//! with excellent concurrency, atomicity, and performance characteristics.
5
6use crate::ManagerPeerEntry;
7use crate::profile as core_profile;
8use eyre::Report;
9use rusqlite::{Connection, OptionalExtension, Transaction, params};
10use serde::{Deserialize, Serialize};
11use std::collections::HashMap;
12use std::path::PathBuf;
13use std::sync::{Arc, Mutex, OnceLock};
14
15fn opt_i64_from_u64(value: Option<u64>) -> Option<i64> {
16    value.and_then(|v| i64::try_from(v).ok())
17}
18
19fn opt_u64_from_i64(value: Option<i64>) -> Option<u64> {
20    value.and_then(|v| u64::try_from(v).ok())
21}
22
23fn u64_from_i64(value: i64) -> u64 {
24    u64::try_from(value).unwrap_or(0)
25}
26
27/// Database connection wrapper with automatic migrations
28pub struct PeerDatabase {
29    conn: Connection,
30}
31
32/// Worker peer entry for database storage
33#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
34pub struct WorkerPeerEntry {
35    pub host: String,
36    pub quic_port: u16,
37    pub tcp_port: u16,
38    pub last_ok: Option<u64>,
39    pub last_fail: Option<u64>,
40    pub last_seen_in_gossip: Option<u64>,
41    pub tls_cert: String,
42    pub tls_fp: String,
43    pub manager_name: String,
44    pub manager_peer: Option<ManagerPeerEntry>,
45}
46
47const SCHEMA_VERSION: i32 = 2;
48
49impl PeerDatabase {
50    /// Create or open a peer database at the given path
51    pub fn open(db_path: PathBuf) -> Result<Self, Report> {
52        // Ensure parent directory exists
53        if let Some(parent) = db_path.parent() {
54            std::fs::create_dir_all(parent)?;
55        }
56
57        let conn = Connection::open(&db_path)?;
58
59        // Configure SQLite for better performance and safety
60        conn.pragma_update(None, "journal_mode", "WAL")?;
61        conn.pragma_update(None, "synchronous", "NORMAL")?;
62        conn.pragma_update(None, "foreign_keys", "ON")?;
63        conn.pragma_update(None, "temp_store", "MEMORY")?;
64
65        let mut db = Self { conn };
66        db.migrate()?;
67        Ok(db)
68    }
69
70    /// Run database migrations to ensure schema is up to date
71    fn migrate(&mut self) -> Result<(), Report> {
72        let tx = self.conn.transaction()?;
73
74        // Create schema_info table if it doesn't exist
75        tx.execute(
76            "CREATE TABLE IF NOT EXISTS schema_info (
77                key TEXT PRIMARY KEY,
78                value INTEGER NOT NULL
79            )",
80            [],
81        )?;
82
83        // Get current schema version
84        let current_version: i32 = tx
85            .query_row(
86                "SELECT value FROM schema_info WHERE key = 'version'",
87                [],
88                |row| row.get(0),
89            )
90            .optional()?
91            .unwrap_or(0);
92
93        if current_version < SCHEMA_VERSION {
94            Self::run_migrations(&tx, current_version)?;
95
96            // Update schema version
97            tx.execute(
98                "INSERT OR REPLACE INTO schema_info (key, value) VALUES ('version', ?)",
99                params![SCHEMA_VERSION],
100            )?;
101        }
102
103        tx.commit()?;
104        Ok(())
105    }
106
107    /// Execute migration scripts based on current version
108    fn run_migrations(tx: &Transaction, from_version: i32) -> Result<(), Report> {
109        if from_version < 1 {
110            // Initial schema creation
111            tx.execute(
112                "CREATE TABLE IF NOT EXISTS manager_peers (
113                    profile TEXT NOT NULL,
114                    manager_id TEXT NOT NULL,
115                    manager_name TEXT NOT NULL,
116                    tenant TEXT NOT NULL,
117                    cluster TEXT NOT NULL,
118                    host TEXT NOT NULL,
119                    tcp_port INTEGER NOT NULL,
120                    quic_port INTEGER NOT NULL,
121                    pub_fp TEXT NOT NULL,
122                    csk_ver INTEGER NOT NULL,
123                    tls_cert TEXT NOT NULL,
124                    tls_fp TEXT NOT NULL,
125                    health_json TEXT,
126                    created_at INTEGER NOT NULL DEFAULT (unixepoch()),
127                    updated_at INTEGER NOT NULL DEFAULT (unixepoch()),
128                    PRIMARY KEY (profile, manager_id, tenant, cluster)
129                )",
130                [],
131            )?;
132
133            tx.execute(
134                "CREATE TABLE IF NOT EXISTS worker_peers (
135                    profile TEXT NOT NULL,
136                    host TEXT NOT NULL,
137                    quic_port INTEGER NOT NULL,
138                    tcp_port INTEGER NOT NULL,
139                    last_ok INTEGER,
140                    last_fail INTEGER,
141                    tls_cert TEXT NOT NULL DEFAULT '',
142                    tls_fp TEXT NOT NULL DEFAULT '',
143                    manager_name TEXT NOT NULL DEFAULT '',
144                    manager_peer_json TEXT,
145                    created_at INTEGER NOT NULL DEFAULT (unixepoch()),
146                    updated_at INTEGER NOT NULL DEFAULT (unixepoch()),
147                    PRIMARY KEY (profile, host, quic_port, tcp_port)
148                )",
149                [],
150            )?;
151
152            tx.execute(
153                "CREATE TABLE IF NOT EXISTS peer_versions (
154                    profile TEXT NOT NULL,
155                    role TEXT NOT NULL,
156                    version INTEGER NOT NULL DEFAULT 0,
157                    updated_at INTEGER NOT NULL DEFAULT (unixepoch()),
158                    PRIMARY KEY (profile, role)
159                )",
160                [],
161            )?;
162
163            // Create indexes for common queries
164            tx.execute(
165                "CREATE INDEX IF NOT EXISTS idx_manager_peers_tenant_cluster
166                 ON manager_peers(profile, tenant, cluster)",
167                [],
168            )?;
169
170            tx.execute(
171                "CREATE INDEX IF NOT EXISTS idx_worker_peers_profile
172                 ON worker_peers(profile)",
173                [],
174            )?;
175
176            // Create triggers to automatically update timestamps
177            tx.execute(
178                "CREATE TRIGGER IF NOT EXISTS update_manager_peers_timestamp
179                 AFTER UPDATE ON manager_peers
180                 BEGIN
181                   UPDATE manager_peers SET updated_at = unixepoch() WHERE rowid = NEW.rowid;
182                 END",
183                [],
184            )?;
185
186            tx.execute(
187                "CREATE TRIGGER IF NOT EXISTS update_worker_peers_timestamp
188                 AFTER UPDATE ON worker_peers
189                 BEGIN
190                   UPDATE worker_peers SET updated_at = unixepoch() WHERE rowid = NEW.rowid;
191                 END",
192                [],
193            )?;
194        }
195
196        if from_version < 2 {
197            // Add last_seen_in_gossip column to existing worker_peers table
198            tx.execute(
199                "ALTER TABLE worker_peers ADD COLUMN last_seen_in_gossip INTEGER",
200                [],
201            )?;
202        }
203
204        Ok(())
205    }
206
207    /// Load all manager peers for a given profile
208    pub fn load_manager_peers(&self, profile: &str) -> Result<Vec<ManagerPeerEntry>, Report> {
209        let mut stmt = self.conn.prepare(
210            "SELECT manager_id, manager_name, tenant, cluster, host, tcp_port, quic_port,
211                    pub_fp, csk_ver, tls_cert, tls_fp, health_json
212             FROM manager_peers
213             WHERE profile = ?
214             ORDER BY manager_id",
215        )?;
216
217        let peer_iter = stmt.query_map([profile], |row| {
218            let health_json: Option<String> = row.get(11)?;
219            let health = match health_json {
220                Some(json) => serde_json::from_str(&json).ok(),
221                None => None,
222            };
223
224            Ok(ManagerPeerEntry {
225                manager_id: row.get(0)?,
226                manager_name: row.get(1)?,
227                tenant: row.get(2)?,
228                cluster: row.get(3)?,
229                host: row.get(4)?,
230                tcp_port: row.get(5)?,
231                quic_port: row.get(6)?,
232                pub_fp: row.get(7)?,
233                csk_ver: row.get(8)?,
234                tls_cert: row.get(9)?,
235                tls_fp: row.get(10)?,
236                health,
237            })
238        })?;
239
240        let mut peers = Vec::new();
241        for peer in peer_iter {
242            peers.push(peer?);
243        }
244
245        tracing::debug!(
246            target: "peer_db",
247            profile = %profile,
248            count = peers.len(),
249            "loaded manager peers from database"
250        );
251
252        Ok(peers)
253    }
254
255    /// Save manager peers for a given profile (replaces all existing peers)
256    pub fn save_manager_peers(
257        &mut self,
258        profile: &str,
259        peers: &[ManagerPeerEntry],
260    ) -> Result<(), Report> {
261        let tx = self.conn.transaction()?;
262
263        // Delete existing peers for this profile
264        tx.execute(
265            "DELETE FROM manager_peers WHERE profile = ?",
266            params![profile],
267        )?;
268
269        // Insert new peers
270        {
271            let mut stmt = tx.prepare(
272                "INSERT INTO manager_peers
273                 (profile, manager_id, manager_name, tenant, cluster, host, tcp_port, quic_port,
274                  pub_fp, csk_ver, tls_cert, tls_fp, health_json)
275                 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
276            )?;
277
278            for peer in peers {
279                let health_json = peer
280                    .health
281                    .as_ref()
282                    .and_then(|h| serde_json::to_string(h).ok());
283
284                stmt.execute(params![
285                    profile,
286                    peer.manager_id,
287                    peer.manager_name,
288                    peer.tenant,
289                    peer.cluster,
290                    peer.host,
291                    peer.tcp_port,
292                    peer.quic_port,
293                    peer.pub_fp,
294                    peer.csk_ver,
295                    peer.tls_cert,
296                    peer.tls_fp,
297                    health_json
298                ])?;
299            }
300        }
301
302        tx.commit()?;
303
304        tracing::debug!(
305            target: "peer_db",
306            profile = %profile,
307            count = peers.len(),
308            "saved manager peers to database"
309        );
310
311        Ok(())
312    }
313
314    /// Add or update a single manager peer
315    pub fn add_manager_peer(
316        &mut self,
317        profile: &str,
318        peer: ManagerPeerEntry,
319    ) -> Result<(), Report> {
320        let tx = self.conn.transaction()?;
321
322        // If we're adding a real peer (not a placeholder), remove any placeholder peers
323        // with the same host/port to prevent having both joining- and real IDs for same node
324        if !peer.manager_id.starts_with("joining-") {
325            tx.execute(
326                "DELETE FROM manager_peers
327                 WHERE profile = ? AND manager_id LIKE 'joining-%'
328                   AND host = ? AND tcp_port = ? AND quic_port = ?
329                   AND tenant = ? AND cluster = ?",
330                params![
331                    profile,
332                    peer.host,
333                    peer.tcp_port,
334                    peer.quic_port,
335                    peer.tenant,
336                    peer.cluster
337                ],
338            )?;
339        }
340
341        let health_json = peer
342            .health
343            .as_ref()
344            .and_then(|h| serde_json::to_string(h).ok());
345
346        // Insert or replace the peer
347        tx.execute(
348            "INSERT OR REPLACE INTO manager_peers
349             (profile, manager_id, manager_name, tenant, cluster, host, tcp_port, quic_port,
350              pub_fp, csk_ver, tls_cert, tls_fp, health_json)
351             VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
352            params![
353                profile,
354                peer.manager_id,
355                peer.manager_name,
356                peer.tenant,
357                peer.cluster,
358                peer.host,
359                peer.tcp_port,
360                peer.quic_port,
361                peer.pub_fp,
362                peer.csk_ver,
363                peer.tls_cert,
364                peer.tls_fp,
365                health_json
366            ],
367        )?;
368
369        tx.commit()?;
370
371        tracing::debug!(
372            target: "peer_db",
373            profile = %profile,
374            manager_id = %peer.manager_id,
375            "added/updated manager peer in database"
376        );
377
378        Ok(())
379    }
380
381    /// Remove a manager peer
382    pub fn remove_manager_peer(
383        &mut self,
384        profile: &str,
385        manager_id: &str,
386        tenant: &str,
387        cluster: &str,
388    ) -> Result<bool, Report> {
389        let rows_affected = self.conn.execute(
390            "DELETE FROM manager_peers
391             WHERE profile = ? AND manager_id = ? AND tenant = ? AND cluster = ?",
392            params![profile, manager_id, tenant, cluster],
393        )?;
394
395        let removed = rows_affected > 0;
396        if removed {
397            tracing::debug!(
398                target: "peer_db",
399                profile = %profile,
400                manager_id = %manager_id,
401                "removed manager peer from database"
402            );
403        }
404
405        Ok(removed)
406    }
407
408    /// Load manager peers for gossip (excludes temporary joining placeholders)
409    pub fn load_manager_peers_for_gossip(
410        &self,
411        profile: &str,
412    ) -> Result<Vec<ManagerPeerEntry>, Report> {
413        let mut stmt = self.conn.prepare(
414            "SELECT manager_id, manager_name, tenant, cluster, host, tcp_port, quic_port,
415                    pub_fp, csk_ver, tls_cert, tls_fp, health_json
416             FROM manager_peers
417             WHERE profile = ? AND manager_id NOT LIKE 'joining-%'
418             ORDER BY manager_id",
419        )?;
420
421        let peer_iter = stmt.query_map([profile], |row| {
422            let health_json: Option<String> = row.get(11)?;
423            let health = match health_json {
424                Some(json) => serde_json::from_str(&json).ok(),
425                None => None,
426            };
427
428            Ok(ManagerPeerEntry {
429                manager_id: row.get(0)?,
430                manager_name: row.get(1)?,
431                tenant: row.get(2)?,
432                cluster: row.get(3)?,
433                host: row.get(4)?,
434                tcp_port: row.get(5)?,
435                quic_port: row.get(6)?,
436                pub_fp: row.get(7)?,
437                csk_ver: row.get(8)?,
438                tls_cert: row.get(9)?,
439                tls_fp: row.get(10)?,
440                health,
441            })
442        })?;
443
444        let mut peers = Vec::new();
445        for peer in peer_iter {
446            peers.push(peer?);
447        }
448
449        Ok(peers)
450    }
451
452    /// Load worker peers for a given profile
453    pub fn load_worker_peers(&self, profile: &str) -> Result<(Vec<WorkerPeerEntry>, u64), Report> {
454        let mut stmt = self.conn.prepare(
455            "SELECT host, quic_port, tcp_port, last_ok, last_fail, last_seen_in_gossip, tls_cert, tls_fp,
456                    manager_name, manager_peer_json
457             FROM worker_peers
458             WHERE profile = ?
459             ORDER BY host, quic_port",
460        )?;
461
462        let peer_iter = stmt.query_map([profile], |row| {
463            let manager_peer_json: Option<String> = row.get(9)?;
464            let manager_peer = match manager_peer_json {
465                Some(json) => serde_json::from_str(&json).ok(),
466                None => None,
467            };
468
469            Ok(WorkerPeerEntry {
470                host: row.get(0)?,
471                quic_port: row.get(1)?,
472                tcp_port: row.get(2)?,
473                last_ok: opt_u64_from_i64(row.get::<_, Option<i64>>(3)?),
474                last_fail: opt_u64_from_i64(row.get::<_, Option<i64>>(4)?),
475                last_seen_in_gossip: opt_u64_from_i64(row.get::<_, Option<i64>>(5)?),
476                tls_cert: row.get(6)?,
477                tls_fp: row.get(7)?,
478                manager_name: row.get(8)?,
479                manager_peer,
480            })
481        })?;
482
483        let mut peers = Vec::new();
484        for peer in peer_iter {
485            peers.push(peer?);
486        }
487
488        let version = self.get_peer_version(profile, "worker")?;
489
490        tracing::debug!(
491            target: "peer_db",
492            profile = %profile,
493            count = peers.len(),
494            version = version,
495            "loaded worker peers from database"
496        );
497
498        Ok((peers, version))
499    }
500
501    /// Save worker peers for a given profile
502    pub fn save_worker_peers(
503        &mut self,
504        profile: &str,
505        peers: &[WorkerPeerEntry],
506        version: u64,
507    ) -> Result<(), Report> {
508        let tx = self.conn.transaction()?;
509
510        // Delete existing peers for this profile
511        tx.execute(
512            "DELETE FROM worker_peers WHERE profile = ?",
513            params![profile],
514        )?;
515
516        // Insert new peers
517        {
518            let mut stmt = tx.prepare(
519                "INSERT INTO worker_peers
520                 (profile, host, quic_port, tcp_port, last_ok, last_fail, last_seen_in_gossip, tls_cert, tls_fp,
521                  manager_name, manager_peer_json)
522                 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
523            )?;
524
525            for peer in peers {
526                let manager_peer_json = peer
527                    .manager_peer
528                    .as_ref()
529                    .and_then(|mp| serde_json::to_string(mp).ok());
530
531                stmt.execute(params![
532                    profile,
533                    peer.host,
534                    peer.quic_port,
535                    peer.tcp_port,
536                    opt_i64_from_u64(peer.last_ok),
537                    opt_i64_from_u64(peer.last_fail),
538                    opt_i64_from_u64(peer.last_seen_in_gossip),
539                    peer.tls_cert,
540                    peer.tls_fp,
541                    peer.manager_name,
542                    manager_peer_json,
543                ])?;
544            }
545        }
546
547        // Update version
548        Self::set_peer_version_tx(&tx, profile, "worker", version)?;
549
550        tx.commit()?;
551
552        tracing::debug!(
553            target: "peer_db",
554            profile = %profile,
555            count = peers.len(),
556            version = version,
557            "saved worker peers to database"
558        );
559
560        Ok(())
561    }
562
563    /// Add a worker peer if it doesn't already exist
564    pub fn add_worker_peer(&mut self, profile: &str, peer: WorkerPeerEntry) -> Result<(), Report> {
565        let manager_peer_json = peer
566            .manager_peer
567            .as_ref()
568            .and_then(|mp| serde_json::to_string(mp).ok());
569
570        let rows_affected = self.conn.execute(
571            "INSERT OR IGNORE INTO worker_peers
572             (profile, host, quic_port, tcp_port, last_ok, last_fail, tls_cert, tls_fp,
573              manager_name, manager_peer_json)
574             VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
575            params![
576                profile,
577                peer.host,
578                peer.quic_port,
579                peer.tcp_port,
580                opt_i64_from_u64(peer.last_ok),
581                opt_i64_from_u64(peer.last_fail),
582                peer.tls_cert,
583                peer.tls_fp,
584                peer.manager_name,
585                manager_peer_json
586            ],
587        )?;
588
589        if rows_affected > 0 {
590            tracing::debug!(
591                target: "peer_db",
592                profile = %profile,
593                host = %peer.host,
594                quic_port = peer.quic_port,
595                "added worker peer to database"
596            );
597        }
598
599        Ok(())
600    }
601
602    /// Get peer version for a given profile and role
603    pub fn get_peer_version(&self, profile: &str, role: &str) -> Result<u64, Report> {
604        let version: i64 = self
605            .conn
606            .query_row(
607                "SELECT version FROM peer_versions WHERE profile = ? AND role = ?",
608                params![profile, role],
609                |row| row.get(0),
610            )
611            .optional()?
612            .unwrap_or(0);
613
614        Ok(u64_from_i64(version))
615    }
616
617    /// Set peer version for a given profile and role
618    pub fn set_peer_version(
619        &mut self,
620        profile: &str,
621        role: &str,
622        version: u64,
623    ) -> Result<(), Report> {
624        let version = i64::try_from(version).unwrap_or(i64::MAX);
625        self.conn.execute(
626            "INSERT OR REPLACE INTO peer_versions (profile, role, version) VALUES (?, ?, ?)",
627            params![profile, role, version],
628        )?;
629
630        Ok(())
631    }
632
633    /// Set peer version within an existing transaction
634    fn set_peer_version_tx(
635        tx: &Transaction,
636        profile: &str,
637        role: &str,
638        version: u64,
639    ) -> Result<(), Report> {
640        let version = i64::try_from(version).unwrap_or(i64::MAX);
641        tx.execute(
642            "INSERT OR REPLACE INTO peer_versions (profile, role, version) VALUES (?, ?, ?)",
643            params![profile, role, version],
644        )?;
645
646        Ok(())
647    }
648
649    /// Increment and return the new peer version for a profile/role
650    pub fn increment_peer_version(&mut self, profile: &str, role: &str) -> Result<u64, Report> {
651        let current_version = self.get_peer_version(profile, role)?;
652        let new_version = current_version + 1;
653        self.set_peer_version(profile, role, new_version)?;
654        Ok(new_version)
655    }
656
657    /// Get database file size in bytes (for monitoring/debugging)
658    pub fn get_db_size(&self) -> Result<u64, Report> {
659        let size: i64 = self.conn.query_row(
660            "SELECT page_count * page_size as size FROM pragma_page_count(), pragma_page_size()",
661            [],
662            |row| row.get(0),
663        )?;
664        Ok(size as u64)
665    }
666
667    /// Optimize database (vacuum and analyze)
668    pub fn optimize(&mut self) -> Result<(), Report> {
669        self.conn.execute("VACUUM", [])?;
670        self.conn.execute("ANALYZE", [])?;
671        Ok(())
672    }
673}
674
675// --- Shared connection pool (role, profile) -> PeerDatabase ---
676
677type DbKey = (String, String); // (role, profile)
678type DbMap = HashMap<DbKey, Arc<Mutex<PeerDatabase>>>;
679static DB_POOL: OnceLock<Arc<Mutex<DbMap>>> = OnceLock::new();
680
681/// Get a shared PeerDatabase for a given role and profile using a global cache.
682///
683/// Connections are created lazily and reused across the process to reduce
684/// SQLite connection churn and ensure consistent WAL settings.
685pub fn get_shared_database(role: &str, profile: &str) -> Result<Arc<Mutex<PeerDatabase>>, Report> {
686    let pool = DB_POOL.get_or_init(|| Arc::new(Mutex::new(HashMap::new())));
687    let mut guard = pool.lock().expect("db pool poisoned");
688    let key: DbKey = (role.to_string(), profile.to_string());
689    if let Some(db) = guard.get(&key) {
690        return Ok(db.clone());
691    }
692
693    // Create new connection at role/profile directory
694    let db_path: PathBuf = core_profile::profile_dir(role, profile).join("peers.db");
695    let db = PeerDatabase::open(db_path)?;
696    let db_arc = Arc::new(Mutex::new(db));
697    guard.insert(key, db_arc.clone());
698    Ok(db_arc)
699}
700
701#[cfg(test)]
702mod tests {
703    use super::*;
704    use crate::HealthMetrics;
705    use tempfile::tempdir;
706
707    fn create_test_manager_peer() -> ManagerPeerEntry {
708        ManagerPeerEntry {
709            manager_id: "test-manager-1".to_string(),
710            manager_name: "Test Manager".to_string(),
711            tenant: "test-tenant".to_string(),
712            cluster: "test-cluster".to_string(),
713            host: "127.0.0.1".to_string(),
714            tcp_port: 8080,
715            quic_port: 8081,
716            pub_fp: "test-fp".to_string(),
717            csk_ver: 1,
718            tls_cert: "test-cert".to_string(),
719            tls_fp: "test-tls-fp".to_string(),
720            health: Some(HealthMetrics {
721                health_score: 0.95,
722                load_percentage: 45.0,
723                max_workers: Some(100),
724                current_workers: 25,
725                avg_cpu: Some(30.5),
726                avg_memory: Some(60.2),
727                last_health_update: 1634567890,
728            }),
729        }
730    }
731
732    fn create_test_worker_peer() -> WorkerPeerEntry {
733        WorkerPeerEntry {
734            host: "192.168.1.100".to_string(),
735            quic_port: 9090,
736            tcp_port: 9091,
737            last_ok: Some(1634567890),
738            last_fail: None,
739            last_seen_in_gossip: None,
740            tls_cert: "worker-cert".to_string(),
741            tls_fp: "worker-fp".to_string(),
742            manager_name: "test-manager".to_string(),
743            manager_peer: None,
744        }
745    }
746
747    #[test]
748    fn test_database_creation_and_migration() {
749        let temp_dir = tempdir().unwrap();
750        let db_path = temp_dir.path().join("test.db");
751
752        let db = PeerDatabase::open(db_path).unwrap();
753
754        // Verify tables were created
755        let table_count: i32 = db.conn.query_row(
756            "SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name IN ('manager_peers', 'worker_peers', 'peer_versions')",
757            [],
758            |row| row.get(0)
759        ).unwrap();
760
761        assert_eq!(table_count, 3);
762    }
763
764    #[test]
765    fn test_manager_peer_operations() {
766        let temp_dir = tempdir().unwrap();
767        let db_path = temp_dir.path().join("test.db");
768        let mut db = PeerDatabase::open(db_path).unwrap();
769
770        let peer = create_test_manager_peer();
771        let profile = "test-profile";
772
773        // Test add peer
774        db.add_manager_peer(profile, peer.clone()).unwrap();
775
776        // Test load peers
777        let loaded_peers = db.load_manager_peers(profile).unwrap();
778        assert_eq!(loaded_peers.len(), 1);
779        assert_eq!(loaded_peers[0].manager_id, peer.manager_id);
780        assert_eq!(loaded_peers[0].health, peer.health);
781
782        // Test remove peer
783        let removed = db
784            .remove_manager_peer(profile, &peer.manager_id, &peer.tenant, &peer.cluster)
785            .unwrap();
786        assert!(removed);
787
788        let loaded_peers = db.load_manager_peers(profile).unwrap();
789        assert_eq!(loaded_peers.len(), 0);
790    }
791
792    #[test]
793    fn test_worker_peer_operations() {
794        let temp_dir = tempdir().unwrap();
795        let db_path = temp_dir.path().join("test.db");
796        let mut db = PeerDatabase::open(db_path).unwrap();
797
798        let peer = create_test_worker_peer();
799        let profile = "test-profile";
800
801        // Test save peers with version
802        db.save_worker_peers(profile, std::slice::from_ref(&peer), 42)
803            .unwrap();
804
805        // Test load peers
806        let (loaded_peers, version) = db.load_worker_peers(profile).unwrap();
807        assert_eq!(loaded_peers.len(), 1);
808        assert_eq!(loaded_peers[0].host, peer.host);
809        assert_eq!(version, 42);
810
811        // Test add peer (should be ignored since it already exists)
812        db.add_worker_peer(profile, peer.clone()).unwrap();
813        let (loaded_peers, _) = db.load_worker_peers(profile).unwrap();
814        assert_eq!(loaded_peers.len(), 1);
815    }
816
817    #[test]
818    fn test_peer_versioning() {
819        let temp_dir = tempdir().unwrap();
820        let db_path = temp_dir.path().join("test.db");
821        let mut db = PeerDatabase::open(db_path).unwrap();
822
823        let profile = "test-profile";
824
825        // Test initial version
826        let version = db.get_peer_version(profile, "manager").unwrap();
827        assert_eq!(version, 0);
828
829        // Test set version
830        db.set_peer_version(profile, "manager", 5).unwrap();
831        let version = db.get_peer_version(profile, "manager").unwrap();
832        assert_eq!(version, 5);
833
834        // Test increment version
835        let new_version = db.increment_peer_version(profile, "manager").unwrap();
836        assert_eq!(new_version, 6);
837    }
838
839    #[test]
840    fn test_joining_peer_replacement() {
841        let temp_dir = tempdir().unwrap();
842        let db_path = temp_dir.path().join("test.db");
843        let mut db = PeerDatabase::open(db_path).unwrap();
844
845        let profile = "test-profile";
846
847        // Add a joining peer
848        let mut joining_peer = create_test_manager_peer();
849        joining_peer.manager_id = "joining-temp-123".to_string();
850        db.add_manager_peer(profile, joining_peer.clone()).unwrap();
851
852        // Add a real peer with same host/port - should remove the joining peer
853        let mut real_peer = create_test_manager_peer();
854        real_peer.manager_id = "real-manager-123".to_string();
855        db.add_manager_peer(profile, real_peer.clone()).unwrap();
856
857        let peers = db.load_manager_peers(profile).unwrap();
858        assert_eq!(peers.len(), 1);
859        assert_eq!(peers[0].manager_id, "real-manager-123");
860
861        // Test gossip filtering
862        let gossip_peers = db.load_manager_peers_for_gossip(profile).unwrap();
863        assert_eq!(gossip_peers.len(), 1);
864        assert_eq!(gossip_peers[0].manager_id, "real-manager-123");
865    }
866}