Skip to main content

gossan_graph/store/
sqlite.rs

1//! SQLite backend with schema versioning and temporal diffing.
2
3use std::path::Path;
4
5use rusqlite::{params, Connection, OptionalExtension, Transaction};
6use chrono::NaiveDateTime;
7
8use crate::schema::{EdgeType, NodeType, SCHEMA_VERSION};
9use crate::store::GraphBackend;
10use crate::{Edge, Node};
11
12/// SQLite-backed graph store.
13pub struct SqliteBackend {
14    conn: Connection,
15}
16
17/// Temporal diff between two scans.
18#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
19pub struct ScanDiff {
20    pub added_targets: Vec<gossan_core::Target>,
21    pub removed_targets: Vec<gossan_core::Target>,
22    pub changed_targets: Vec<gossan_core::Target>,
23    pub added_findings: Vec<secfinding::Finding>,
24    pub removed_findings: Vec<secfinding::Finding>,
25    pub changed_findings: Vec<secfinding::Finding>,
26}
27
28impl SqliteBackend {
29    /// Convert milliseconds timestamp to SQLite datetime string.
30    fn ms_to_datetime(ms: u64) -> String {
31        let seconds = (ms / 1000) as i64;
32        match chrono::DateTime::from_timestamp(seconds, 0) {
33            Some(dt) => dt.format("%Y-%m-%d %H:%M:%S").to_string(),
34            None => "1970-01-01 00:00:00".to_string(),
35        }
36    }
37
38    /// Convert SQLite datetime string to milliseconds timestamp.
39    fn datetime_to_ms(dt_str: &str) -> u64 {
40        match NaiveDateTime::parse_from_str(dt_str, "%Y-%m-%d %H:%M:%S") {
41            Ok(dt) => dt.and_utc().timestamp() as u64 * 1000,
42            Err(_) => 0,
43        }
44    }
45
46    /// Open or create a SQLite graph database.
47    pub fn open<P: AsRef<Path>>(path: P) -> Result<Self, SqliteError> {
48        let conn = Connection::open(path)?;
49        conn.execute_batch(
50            "PRAGMA journal_mode = WAL;
51             PRAGMA synchronous = NORMAL;
52             PRAGMA foreign_keys = ON;
53             PRAGMA busy_timeout = 5000;",
54        )?;
55        let mut backend = Self { conn };
56        backend.init_schema()?;
57        Ok(backend)
58    }
59
60    /// Open an in-memory database for testing.
61    #[cfg(test)]
62    pub fn open_in_memory() -> Result<Self, SqliteError> {
63        let conn = Connection::open_in_memory()?;
64        let mut backend = Self { conn };
65        backend.init_schema()?;
66        Ok(backend)
67    }
68
69    fn init_schema(&mut self) -> Result<(), SqliteError> {
70        self.conn.execute_batch(
71            "CREATE TABLE IF NOT EXISTS schema_version (
72                version INTEGER PRIMARY KEY,
73                updated_at DATETIME DEFAULT CURRENT_TIMESTAMP
74            );
75            CREATE TABLE IF NOT EXISTS targets (
76                id TEXT PRIMARY KEY,
77                kind TEXT NOT NULL,
78                label TEXT NOT NULL,
79                data TEXT,
80                first_seen DATETIME DEFAULT CURRENT_TIMESTAMP,
81                last_seen DATETIME DEFAULT CURRENT_TIMESTAMP
82            );
83            CREATE TABLE IF NOT EXISTS findings (
84                id TEXT PRIMARY KEY,
85                kind TEXT NOT NULL DEFAULT 'finding',
86                label TEXT NOT NULL,
87                data TEXT,
88                first_seen DATETIME DEFAULT CURRENT_TIMESTAMP,
89                last_seen DATETIME DEFAULT CURRENT_TIMESTAMP
90            );
91            CREATE TABLE IF NOT EXISTS relationships (
92                source_id TEXT NOT NULL,
93                target_id TEXT NOT NULL,
94                rel_type TEXT NOT NULL,
95                data TEXT,
96                first_seen DATETIME DEFAULT CURRENT_TIMESTAMP,
97                last_seen DATETIME DEFAULT CURRENT_TIMESTAMP,
98                PRIMARY KEY (source_id, target_id, rel_type)
99            );
100            CREATE INDEX IF NOT EXISTS idx_relationships_source ON relationships(source_id);
101            CREATE INDEX IF NOT EXISTS idx_relationships_target ON relationships(target_id);
102            CREATE INDEX IF NOT EXISTS idx_targets_kind ON targets(kind);
103            CREATE INDEX IF NOT EXISTS idx_findings_kind ON findings(kind);
104            ",
105        )?;
106
107        let current: Option<i64> = self
108            .conn
109            .query_row(
110                "SELECT version FROM schema_version ORDER BY version DESC LIMIT 1",
111                [],
112                |row| row.get(0),
113            )
114            .optional()?;
115
116        match current {
117            None => {
118                self.conn.execute(
119                    "INSERT INTO schema_version (version) VALUES (?1)",
120                    params![SCHEMA_VERSION],
121                )?;
122            }
123            Some(v) => {
124                let v = v as u32;
125                if v > SCHEMA_VERSION {
126                    return Err(SqliteError::Schema(
127                        crate::schema::SchemaError::UnsupportedVersion {
128                            found: v,
129                            max_supported: SCHEMA_VERSION,
130                        }
131                        .to_string(),
132                    ));
133                }
134                // Run migrations here when SCHEMA_VERSION is bumped.
135                if v < SCHEMA_VERSION {
136                    self.migrate(v, SCHEMA_VERSION)?;
137                }
138            }
139        }
140
141        Ok(())
142    }
143
144    fn migrate(&mut self, from: u32, to: u32) -> Result<(), SqliteError> {
145        let tx = self.conn.transaction()?;
146        // Placeholder for future migrations.
147        // Example:
148        // if from < 2 {
149        //     tx.execute("ALTER TABLE nodes ADD COLUMN new_col TEXT", [])?;
150        // }
151        tx.execute(
152            "INSERT INTO schema_version (version) VALUES (?1)",
153            params![to],
154        )?;
155        tx.commit()?;
156        tracing::info!(from, to, "graph schema migrated");
157        Ok(())
158    }
159
160    /// Persist a scan of targets and findings, inferring edges.
161    pub fn persist_scan(
162        &mut self,
163        targets: &[gossan_core::Target],
164        findings: &[secfinding::Finding],
165    ) -> Result<(), SqliteError> {
166        let tx = self.conn.transaction()?;
167
168        for target in targets {
169            let node = target_to_node(target);
170            Self::upsert_node(&tx, &node)?;
171            Self::insert_inferred_target_edges(&tx, target)?;
172        }
173
174        for finding in findings {
175            let node = finding_to_node(finding);
176            Self::upsert_node(&tx, &node)?;
177
178            let target_id = target_id_from_finding(finding)?;
179            // A finding's target row must exist before we can hang an
180            // edge on it. Callers may legitimately persist findings
181            // without first persisting the matching Target (e.g. an
182            // out-of-band scanner that only emits findings). Insert a
183            // stub Target node — INSERT OR IGNORE means real Target
184            // payloads from a same-transaction targets[] entry win.
185            Self::insert_stub_target_for_finding(&tx, &target_id, finding)?;
186            let edge = Edge::new(&target_id, &node.id, EdgeType::HasFinding);
187            Self::upsert_edge(&tx, &edge)?;
188        }
189
190        tx.commit()?;
191        Ok(())
192    }
193
194    /// Compute temporal diff.
195    pub fn compute_diff(
196        &self,
197        targets: &[gossan_core::Target],
198        findings: &[secfinding::Finding],
199        removed_threshold: std::time::Duration,
200    ) -> Result<ScanDiff, SqliteError> {
201        let mut diff = ScanDiff {
202            added_targets: Vec::new(),
203            removed_targets: Vec::new(),
204            changed_targets: Vec::new(),
205            added_findings: Vec::new(),
206            removed_findings: Vec::new(),
207            changed_findings: Vec::new(),
208        };
209
210        for target in targets {
211            let id = target_id(target);
212            let existing: Option<String> = self
213                .conn
214                .query_row(
215                    "SELECT data FROM targets WHERE id = ?1",
216                    params![id],
217                    |row| row.get(0),
218                )
219                .optional()?;
220            // Compare structurally, not by raw string. The stored `data`
221            // column was written by `Node::with_payload` (which goes
222            // target → serde_json::Value → Value::to_string) while the
223            // diff side serialises target directly. Both should be
224            // semantically identical, but key ordering / whitespace /
225            // numeric formatting are not contractually byte-equal across
226            // those two paths. Decoding both sides to a comparable shape
227            // is the only way to ask "did this target actually change?".
228            match existing {
229                None => diff.added_targets.push(target.clone()),
230                Some(old) => {
231                    let stored_val: serde_json::Value =
232                        serde_json::from_str(&old).unwrap_or(serde_json::Value::Null);
233                    let new_val = serde_json::to_value(target)?;
234                    if stored_val != new_val {
235                        diff.changed_targets.push(target.clone());
236                    }
237                }
238            }
239        }
240
241        for finding in findings {
242            let id = finding_id(finding);
243            let existing: Option<String> = self
244                .conn
245                .query_row(
246                    "SELECT data FROM findings WHERE id = ?1",
247                    params![id],
248                    |row| row.get(0),
249                )
250                .optional()?;
251            match existing {
252                None => diff.added_findings.push(finding.clone()),
253                Some(old) => {
254                    let stored_val: serde_json::Value =
255                        serde_json::from_str(&old).unwrap_or(serde_json::Value::Null);
256                    let new_val = serde_json::to_value(finding)?;
257                    if stored_val != new_val {
258                        diff.changed_findings.push(finding.clone());
259                    }
260                }
261            }
262        }
263
264        // Clamp threshold to SQLite's practical limit (~100 years in seconds)
265        let threshold_secs = removed_threshold.as_secs().min(3_153_600_000u64) as i64;
266        let threshold_datetime = format!("-{} seconds", threshold_secs);
267
268        let mut stmt = self.conn.prepare(
269            "SELECT data FROM targets 
270             WHERE kind IN ('domain','host','service','web','network','repository','package')
271               AND last_seen <= datetime('now', ?1)",
272        )?;
273        let rows = stmt.query_map(params![threshold_datetime], |row| {
274            let data: String = row.get(0)?;
275            serde_json::from_str::<gossan_core::Target>(&data)
276                .map_err(|e| rusqlite::Error::FromSqlConversionFailure(
277                    0,
278                    rusqlite::types::Type::Text,
279                    Box::new(e),
280                ))
281        })?;
282        for r in rows {
283            diff.removed_targets.push(r?);
284        }
285
286        let mut stmt = self.conn.prepare(
287            "SELECT data FROM findings 
288             WHERE last_seen <= datetime('now', ?1)",
289        )?;
290        let rows = stmt.query_map(params![threshold_datetime], |row| {
291            let data: String = row.get(0)?;
292            serde_json::from_str::<secfinding::Finding>(&data).map_err(|e| {
293                rusqlite::Error::FromSqlConversionFailure(
294                    0,
295                    rusqlite::types::Type::Text,
296                    Box::new(e),
297                )
298            })
299        })?;
300        for r in rows {
301            diff.removed_findings.push(r?);
302        }
303
304        Ok(diff)
305    }
306
307    fn upsert_node(tx: &Transaction, node: &Node) -> Result<(), SqliteError> {
308        let data = node
309            .payload
310            .as_ref()
311            .map(|p| p.to_string())
312            .unwrap_or_default();
313        let first_seen = Self::ms_to_datetime(node.first_seen_ms);
314        let last_seen = Self::ms_to_datetime(node.last_seen_ms);
315        
316        // Determine table based on node type
317        let table = if node.kind == NodeType::Finding {
318            "findings"
319        } else {
320            "targets"
321        };
322        
323        let insert_query = format!(
324            "INSERT OR IGNORE INTO {} (id, kind, label, data, first_seen, last_seen) 
325             VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
326            table
327        );
328        tx.execute(
329            &insert_query,
330            params![
331                node.id,
332                node.kind.to_string(),
333                node.label,
334                data,
335                first_seen,
336                last_seen
337            ],
338        )?;
339        
340        let update_query = format!(
341            "UPDATE {} SET last_seen = ?2, data = ?3 WHERE id = ?1",
342            table
343        );
344        tx.execute(
345            &update_query,
346            params![node.id, last_seen, data],
347        )?;
348        Ok(())
349    }
350
351    fn upsert_edge(tx: &Transaction, edge: &Edge) -> Result<(), SqliteError> {
352        let data = edge
353            .payload
354            .as_ref()
355            .map(|p| p.to_string())
356            .unwrap_or_default();
357        let first_seen = Self::ms_to_datetime(edge.first_seen_ms);
358        let last_seen = Self::ms_to_datetime(edge.last_seen_ms);
359        
360        tx.execute(
361            "INSERT OR IGNORE INTO relationships (source_id, target_id, rel_type, data, first_seen, last_seen)
362             VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
363            params![
364                edge.source_id,
365                edge.target_id,
366                edge.kind.to_string(),
367                data,
368                first_seen,
369                last_seen
370            ],
371        )?;
372        tx.execute(
373            "UPDATE relationships SET last_seen = ?4, data = ?5
374             WHERE source_id = ?1 AND target_id = ?2 AND rel_type = ?3",
375            params![edge.source_id, edge.target_id, edge.kind.to_string(), last_seen, data],
376        )?;
377        Ok(())
378    }
379
380    fn insert_stub_target_for_finding(
381        tx: &Transaction,
382        target_id: &str,
383        finding: &secfinding::Finding,
384    ) -> Result<(), SqliteError> {
385        // Derive NodeType from the target_id's prefix (set by
386        // target_id_from_finding). Default to Endpoint for anything
387        // not in the known set so we don't lose the row.
388        let kind = match target_id.split_once(':').map(|(p, _)| p) {
389            Some("domain") => NodeType::Domain,
390            Some("host") => NodeType::Ip,
391            Some("service") => NodeType::Service,
392            Some("web") => NodeType::Endpoint,
393            _ => NodeType::Endpoint,
394        };
395        let label = finding.target().to_string();
396        let now = std::time::SystemTime::now()
397            .duration_since(std::time::UNIX_EPOCH)
398            .unwrap_or_default()
399            .as_millis() as u64;
400        let now_dt = Self::ms_to_datetime(now);
401        // INSERT OR IGNORE only — never clobber a real Target row that
402        // was upserted earlier in the same transaction with a real
403        // payload.
404        tx.execute(
405            "INSERT OR IGNORE INTO targets (id, kind, label, data, first_seen, last_seen)
406             VALUES (?1, ?2, ?3, '', ?4, ?5)",
407            params![target_id, kind.to_string(), label, now_dt, now_dt],
408        )?;
409        Ok(())
410    }
411
412    fn insert_inferred_target_edges(
413        tx: &Transaction,
414        target: &gossan_core::Target,
415    ) -> Result<(), SqliteError> {
416        match target {
417            gossan_core::Target::Host(h) => {
418                if let Some(domain) = &h.domain {
419                    let e = Edge::new(
420                        format!("domain:{domain}"),
421                        format!("host:{}", h.ip),
422                        EdgeType::ResolvesTo,
423                    );
424                    Self::upsert_edge(tx, &e)?;
425                }
426            }
427            gossan_core::Target::Service(s) => {
428                let e = Edge::new(
429                    format!("host:{}", s.host.ip),
430                    format!("service:{}:{}", s.host.ip, s.port),
431                    EdgeType::HasService,
432                );
433                Self::upsert_edge(tx, &e)?;
434                if let Some(domain) = &s.host.domain {
435                    let e = Edge::new(
436                        format!("domain:{domain}"),
437                        format!("service:{}:{}", s.host.ip, s.port),
438                        EdgeType::HasService,
439                    );
440                    Self::upsert_edge(tx, &e)?;
441                }
442            }
443            gossan_core::Target::Web(w) => {
444                let e = Edge::new(
445                    format!("service:{}:{}", w.service.host.ip, w.service.port),
446                    format!("web:{}", w.url),
447                    EdgeType::Exposes,
448                );
449                Self::upsert_edge(tx, &e)?;
450            }
451            _ => {}
452        }
453        Ok(())
454    }
455
456    /// Raw SQL escape hatch for advanced queries.
457    pub fn conn(&self) -> &Connection {
458        &self.conn
459    }
460}
461
462impl GraphBackend for SqliteBackend {
463    type Error = SqliteError;
464
465    fn init(&mut self) -> Result<(), Self::Error> {
466        self.init_schema()?;
467        Ok(())
468    }
469
470    fn write_nodes(&mut self, nodes: &[Node]) -> Result<(), Self::Error> {
471        let tx = self.conn.transaction()?;
472        for n in nodes {
473            Self::upsert_node(&tx, n)?;
474        }
475        tx.commit()?;
476        Ok(())
477    }
478
479    fn write_edges(&mut self, edges: &[Edge]) -> Result<(), Self::Error> {
480        let tx = self.conn.transaction()?;
481        for e in edges {
482            Self::upsert_edge(&tx, e)?;
483        }
484        tx.commit()?;
485        Ok(())
486    }
487
488    fn read_nodes(&self) -> Result<Vec<Node>, Self::Error> {
489        let mut nodes = Vec::new();
490        
491        // Read from targets table
492        let mut stmt = self.conn.prepare(
493            "SELECT id, kind, label, data, first_seen, last_seen FROM targets",
494        )?;
495        let target_rows = stmt.query_map([], |row| {
496            let kind_str: String = row.get(1)?;
497            let data_str: String = row.get(3)?;
498            let first_seen_str: String = row.get(4)?;
499            let last_seen_str: String = row.get(5)?;
500            Ok(Node {
501                id: row.get(0)?,
502                kind: parse_node_type(&kind_str).unwrap_or(NodeType::Domain),
503                label: row.get(2)?,
504                payload: if data_str.is_empty() {
505                    None
506                } else {
507                    serde_json::from_str(&data_str).ok()
508                },
509                first_seen_ms: Self::datetime_to_ms(&first_seen_str),
510                last_seen_ms: Self::datetime_to_ms(&last_seen_str),
511            })
512        })?;
513        for row in target_rows {
514            nodes.push(row?);
515        }
516        
517        // Read from findings table
518        let mut stmt = self.conn.prepare(
519            "SELECT id, kind, label, data, first_seen, last_seen FROM findings",
520        )?;
521        let finding_rows = stmt.query_map([], |row| {
522            let kind_str: String = row.get(1)?;
523            let data_str: String = row.get(3)?;
524            let first_seen_str: String = row.get(4)?;
525            let last_seen_str: String = row.get(5)?;
526            Ok(Node {
527                id: row.get(0)?,
528                kind: parse_node_type(&kind_str).unwrap_or(NodeType::Finding),
529                label: row.get(2)?,
530                payload: if data_str.is_empty() {
531                    None
532                } else {
533                    serde_json::from_str(&data_str).ok()
534                },
535                first_seen_ms: Self::datetime_to_ms(&first_seen_str),
536                last_seen_ms: Self::datetime_to_ms(&last_seen_str),
537            })
538        })?;
539        for row in finding_rows {
540            nodes.push(row?);
541        }
542        
543        Ok(nodes)
544    }
545
546    fn read_edges(&self) -> Result<Vec<Edge>, Self::Error> {
547        let mut stmt = self.conn.prepare(
548            "SELECT source_id, target_id, rel_type, data, first_seen, last_seen FROM relationships",
549        )?;
550        let rows = stmt.query_map([], |row| {
551            let kind_str: String = row.get(2)?;
552            let data_str: String = row.get(3)?;
553            let first_seen_str: String = row.get(4)?;
554            let last_seen_str: String = row.get(5)?;
555            Ok(Edge {
556                source_id: row.get(0)?,
557                target_id: row.get(1)?,
558                kind: parse_edge_type(&kind_str).unwrap_or(EdgeType::HasFinding),
559                payload: if data_str.is_empty() {
560                    None
561                } else {
562                    serde_json::from_str(&data_str).ok()
563                },
564                first_seen_ms: Self::datetime_to_ms(&first_seen_str),
565                last_seen_ms: Self::datetime_to_ms(&last_seen_str),
566            })
567        })?;
568        rows.collect::<Result<Vec<_>, _>>()
569            .map_err(Into::into)
570    }
571
572    fn find_nodes_by_type(&self, kind: NodeType) -> Result<Vec<Node>, Self::Error> {
573        let kind_str = kind.to_string();
574        let table = if kind == NodeType::Finding {
575            "findings"
576        } else {
577            "targets"
578        };
579        
580        let query = format!(
581            "SELECT id, kind, label, data, first_seen, last_seen FROM {} WHERE kind = ?1",
582            table
583        );
584        let mut stmt = self.conn.prepare(&query)?;
585        let rows = stmt.query_map(params![kind_str], |row| {
586            let data_str: String = row.get(3)?;
587            let first_seen_str: String = row.get(4)?;
588            let last_seen_str: String = row.get(5)?;
589            Ok(Node {
590                id: row.get(0)?,
591                kind: kind.clone(),
592                label: row.get(2)?,
593                payload: if data_str.is_empty() {
594                    None
595                } else {
596                    serde_json::from_str(&data_str).ok()
597                },
598                first_seen_ms: Self::datetime_to_ms(&first_seen_str),
599                last_seen_ms: Self::datetime_to_ms(&last_seen_str),
600            })
601        })?;
602        rows.collect::<Result<Vec<_>, _>>()
603            .map_err(Into::into)
604    }
605
606    fn neighbors(
607        &self,
608        node_id: &str,
609        edge_type: Option<EdgeType>,
610    ) -> Result<Vec<Edge>, Self::Error> {
611        let node_id = node_id.to_string();
612        let mut stmt = match edge_type {
613            Some(ref et) => self.conn.prepare(
614                "SELECT source_id, target_id, rel_type, data, first_seen, last_seen
615                 FROM relationships WHERE source_id = ?1 AND rel_type = ?2",
616            )?,
617            None => self.conn.prepare(
618                "SELECT source_id, target_id, rel_type, data, first_seen, last_seen
619                 FROM relationships WHERE source_id = ?1",
620            )?,
621        };
622        let map_row = |row: &rusqlite::Row<'_>| -> Result<Edge, rusqlite::Error> {
623            let kind_str: String = row.get(2)?;
624            let data_str: String = row.get(3)?;
625            let first_seen_str: String = row.get(4)?;
626            let last_seen_str: String = row.get(5)?;
627            Ok(Edge {
628                source_id: row.get(0)?,
629                target_id: row.get(1)?,
630                kind: parse_edge_type(&kind_str).unwrap_or(EdgeType::HasFinding),
631                payload: if data_str.is_empty() {
632                    None
633                } else {
634                    serde_json::from_str(&data_str).ok()
635                },
636                first_seen_ms: Self::datetime_to_ms(&first_seen_str),
637                last_seen_ms: Self::datetime_to_ms(&last_seen_str),
638            })
639        };
640        let rows = match edge_type {
641            Some(et) => stmt.query_map(params![node_id, et.to_string()], map_row)?,
642            None => stmt.query_map(params![node_id], map_row)?,
643        };
644        rows.collect::<Result<Vec<_>, _>>()
645            .map_err(Into::into)
646    }
647
648    fn clear(&mut self) -> Result<(), Self::Error> {
649        self.conn.execute("DELETE FROM relationships", [])?;
650        self.conn.execute("DELETE FROM findings", [])?;
651        self.conn.execute("DELETE FROM targets", [])?;
652        Ok(())
653    }
654}
655
656/// Error type for SQLite backend operations.
657#[derive(Debug, thiserror::Error)]
658pub enum SqliteError {
659    #[error("SQLite error: {0}")]
660    Sqlite(#[from] rusqlite::Error),
661    #[error("JSON error: {0}")]
662    Json(#[from] serde_json::Error),
663    #[error("Schema error: {0}")]
664    Schema(String),
665}
666
667/// Compute the deterministic node ID for a `Target`. Used internally
668/// by the SQLite backend (as the `target_id` column on `edges` /
669/// `findings`) and by external callers that want to round-trip
670/// Target identity through a graph store. Promoted to `pub` so the
671/// legendary unit test can pin the ID format
672/// (`domain:<host>` / `host:<ip>` / `service:<ip>:<port>` / etc.)
673/// without depending on the `target_id_from_finding` flavour, which
674/// works on findings instead of targets.
675pub fn target_id(target: &gossan_core::Target) -> String {
676    match target {
677        gossan_core::Target::Domain(d) => format!("domain:{}", d.domain),
678        gossan_core::Target::Host(h) => format!("host:{}", h.ip),
679        gossan_core::Target::Service(s) => format!("service:{}:{}", s.host.ip, s.port),
680        gossan_core::Target::Web(w) => format!("web:{}", w.url),
681        gossan_core::Target::Network(n) => format!("network:{}", n.cidr),
682        gossan_core::Target::Repository(r) => format!("repo:{}", r.url),
683        gossan_core::Target::InternalPackage(p) => format!("pkg:{}", p.name),
684        _ => {
685            let data = serde_json::to_string(target).unwrap_or_default();
686            format!("unknown:{}", &data[..data.len().min(120)])
687        }
688    }
689}
690
691fn target_to_node(target: &gossan_core::Target) -> Node {
692    let id = target_id(target);
693    let (kind, label) = match target {
694        gossan_core::Target::Domain(d) => (NodeType::Domain, d.domain.clone()),
695        gossan_core::Target::Host(h) => (NodeType::Ip, h.ip.to_string()),
696        gossan_core::Target::Service(s) => {
697            (NodeType::Service, format!("{}:{}", s.host.ip, s.port))
698        }
699        gossan_core::Target::Web(w) => (NodeType::Endpoint, w.url.to_string()),
700        gossan_core::Target::Network(n) => (NodeType::Ip, n.cidr.clone()),
701        gossan_core::Target::Repository(r) => (NodeType::Endpoint, r.url.to_string()),
702        gossan_core::Target::InternalPackage(p) => (NodeType::Endpoint, p.name.clone()),
703        _ => (NodeType::Endpoint, id.clone()),
704    };
705    Node::new(id, kind, label).with_payload(target)
706}
707
708fn finding_id(finding: &secfinding::Finding) -> String {
709    let namespace = uuid::Uuid::NAMESPACE_OID;
710    let content = format!(
711        "{}:{}:{:?}:{}",
712        finding.target(),
713        finding.title(),
714        finding.severity(),
715        finding.detail()
716    );
717    let id = uuid::Uuid::new_v5(&namespace, content.as_bytes());
718    format!("finding:{id}")
719}
720
721fn finding_to_node(finding: &secfinding::Finding) -> Node {
722    let id = finding_id(finding);
723    Node::new(id, NodeType::Finding, finding.title().to_string()).with_payload(finding)
724}
725
726/// Derive a target node id from a finding target string.
727///
728/// # Errors
729///
730/// Returns an error if the target string cannot be parsed into a known shape.
731pub fn target_id_from_finding(finding: &secfinding::Finding) -> Result<String, SqliteError> {
732    let t = finding.target();
733
734    // Try URL first — but only treat it as a Web target if the parser
735    // actually saw a recognized HTTP-family scheme. `url::Url::parse`
736    // is happy to interpret `"example.com:443"` as `scheme=example.com,
737    // path=443`, which would misclassify a bare host:port pair as a
738    // Web URL.
739    if let Ok(url) = url::Url::parse(t) {
740        if matches!(url.scheme(), "http" | "https" | "ws" | "wss" | "ftp") {
741            return Ok(format!("web:{}", url));
742        }
743    }
744
745    // Try IP address (IPv4 and IPv6)
746    if t.parse::<std::net::IpAddr>().is_ok() {
747        return Ok(format!("host:{t}"));
748    }
749
750    // Try bracketed IPv6
751    if t.starts_with('[') && t.contains("]:") {
752        if let Some(idx) = t.find(']') {
753            let ip_part = &t[1..idx];
754            if ip_part.parse::<std::net::IpAddr>().is_ok() {
755                return Ok(format!("service:{t}"));
756            }
757        }
758    }
759
760    // host:port or ip:port — but avoid misclassifying domains like example.com:443
761    if let Some((host, port)) = t.rsplit_once(':') {
762        if port.parse::<u16>().is_ok() {
763            if host.parse::<std::net::IpAddr>().is_ok() {
764                return Ok(format!("service:{t}"));
765            }
766            // If it looks like a bare IPv6 without brackets, reject rather than guess.
767            if host.contains(':') {
768                return Err(SqliteError::Schema(format!(
769                    "ambiguous IPv6 service target without brackets: {t}"
770                )));
771            }
772        }
773    }
774
775    // Default: domain
776    Ok(format!("domain:{t}"))
777}
778
779fn parse_node_type(s: &str) -> Option<NodeType> {
780    match s {
781        "domain" => Some(NodeType::Domain),
782        "subdomain" => Some(NodeType::Subdomain),
783        "ip" => Some(NodeType::Ip),
784        "port" => Some(NodeType::Port),
785        "service" => Some(NodeType::Service),
786        "tech" => Some(NodeType::Tech),
787        "endpoint" => Some(NodeType::Endpoint),
788        "secret" => Some(NodeType::Secret),
789        "cloud" => Some(NodeType::Cloud),
790        "finding" => Some(NodeType::Finding),
791        _ => None,
792    }
793}
794
795fn parse_edge_type(s: &str) -> Option<EdgeType> {
796    match s {
797        "RESOLVES_TO" => Some(EdgeType::ResolvesTo),
798        "HOSTS" => Some(EdgeType::Hosts),
799        "RUNS" => Some(EdgeType::Runs),
800        "EXPOSES" => Some(EdgeType::Exposes),
801        "LEAKS" => Some(EdgeType::Leaks),
802        "MISCONFIGURED" => Some(EdgeType::Misconfigured),
803        "HAS_FINDING" => Some(EdgeType::HasFinding),
804        "HAS_SERVICE" => Some(EdgeType::HasService),
805        _ => None,
806    }
807}
808
809#[cfg(test)]
810mod tests {
811    use super::*;
812
813    #[test]
814    fn target_id_from_finding_url() {
815        let f = secfinding::Finding::new(
816            "s",
817            "https://example.com/path",
818            secfinding::Severity::Info,
819            "t",
820            "",
821        )
822        .unwrap();
823        assert_eq!(
824            target_id_from_finding(&f).unwrap(),
825            "web:https://example.com/path"
826        );
827    }
828
829    #[test]
830    fn target_id_from_finding_ipv4() {
831        let f = secfinding::Finding::new("s", "1.2.3.4", secfinding::Severity::Info, "t", "")
832            .unwrap();
833        assert_eq!(target_id_from_finding(&f).unwrap(), "host:1.2.3.4");
834    }
835
836    #[test]
837    fn target_id_from_finding_ipv6() {
838        let f =
839            secfinding::Finding::new("s", "::1", secfinding::Severity::Info, "t", "").unwrap();
840        assert_eq!(target_id_from_finding(&f).unwrap(), "host:::1");
841    }
842
843    #[test]
844    fn target_id_from_finding_service() {
845        let f =
846            secfinding::Finding::new("s", "1.2.3.4:443", secfinding::Severity::Info, "t", "")
847                .unwrap();
848        assert_eq!(target_id_from_finding(&f).unwrap(), "service:1.2.3.4:443");
849    }
850
851    #[test]
852    fn target_id_from_finding_domain_with_port() {
853        let f = secfinding::Finding::new(
854            "s",
855            "example.com:443",
856            secfinding::Severity::Info,
857            "t",
858            "",
859        )
860        .unwrap();
861        // Domain with port but no scheme falls through to domain.
862        assert_eq!(target_id_from_finding(&f).unwrap(), "domain:example.com:443");
863    }
864
865    #[test]
866    fn sqlite_roundtrip() {
867        let mut backend = SqliteBackend::open_in_memory().unwrap();
868        let node = Node::new("n1", NodeType::Domain, "example.com");
869        backend.write_nodes(&[node]).unwrap();
870
871        let edge = Edge::new("n1", "n2", EdgeType::ResolvesTo);
872        backend.write_edges(&[edge]).unwrap();
873
874        let nodes = backend.read_nodes().unwrap();
875        assert_eq!(nodes.len(), 1);
876        assert_eq!(nodes[0].id, "n1");
877
878        let edges = backend.read_edges().unwrap();
879        assert_eq!(edges.len(), 1);
880        assert_eq!(edges[0].kind, EdgeType::ResolvesTo);
881    }
882
883    #[test]
884    fn schema_version_tracked() {
885        let backend = SqliteBackend::open_in_memory().unwrap();
886        let v: i64 = backend
887            .conn()
888            .query_row(
889                "SELECT version FROM schema_version ORDER BY version DESC LIMIT 1",
890                [],
891                |row| row.get(0),
892            )
893            .unwrap();
894        assert_eq!(v, i64::from(SCHEMA_VERSION));
895    }
896}