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 chrono::NaiveDateTime;
6use rusqlite::{params, Connection, OptionalExtension, Transaction};
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).map_err(|e| {
276                rusqlite::Error::FromSqlConversionFailure(
277                    0,
278                    rusqlite::types::Type::Text,
279                    Box::new(e),
280                )
281            })
282        })?;
283        for r in rows {
284            diff.removed_targets.push(r?);
285        }
286
287        let mut stmt = self.conn.prepare(
288            "SELECT data FROM findings 
289             WHERE last_seen <= datetime('now', ?1)",
290        )?;
291        let rows = stmt.query_map(params![threshold_datetime], |row| {
292            let data: String = row.get(0)?;
293            serde_json::from_str::<secfinding::Finding>(&data).map_err(|e| {
294                rusqlite::Error::FromSqlConversionFailure(
295                    0,
296                    rusqlite::types::Type::Text,
297                    Box::new(e),
298                )
299            })
300        })?;
301        for r in rows {
302            diff.removed_findings.push(r?);
303        }
304
305        Ok(diff)
306    }
307
308    fn upsert_node(tx: &Transaction, node: &Node) -> Result<(), SqliteError> {
309        let data = node
310            .payload
311            .as_ref()
312            .map(|p| p.to_string())
313            .unwrap_or_default();
314        let first_seen = Self::ms_to_datetime(node.first_seen_ms);
315        let last_seen = Self::ms_to_datetime(node.last_seen_ms);
316
317        // Determine table based on node type
318        let table = if node.kind == NodeType::Finding {
319            "findings"
320        } else {
321            "targets"
322        };
323
324        let insert_query = format!(
325            "INSERT OR IGNORE INTO {} (id, kind, label, data, first_seen, last_seen) 
326             VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
327            table
328        );
329        tx.execute(
330            &insert_query,
331            params![
332                node.id,
333                node.kind.to_string(),
334                node.label,
335                data,
336                first_seen,
337                last_seen
338            ],
339        )?;
340
341        let update_query = format!(
342            "UPDATE {} SET last_seen = ?2, data = ?3 WHERE id = ?1",
343            table
344        );
345        tx.execute(&update_query, params![node.id, last_seen, data])?;
346        Ok(())
347    }
348
349    fn upsert_edge(tx: &Transaction, edge: &Edge) -> Result<(), SqliteError> {
350        let data = edge
351            .payload
352            .as_ref()
353            .map(|p| p.to_string())
354            .unwrap_or_default();
355        let first_seen = Self::ms_to_datetime(edge.first_seen_ms);
356        let last_seen = Self::ms_to_datetime(edge.last_seen_ms);
357
358        tx.execute(
359            "INSERT OR IGNORE INTO relationships (source_id, target_id, rel_type, data, first_seen, last_seen)
360             VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
361            params![
362                edge.source_id,
363                edge.target_id,
364                edge.kind.to_string(),
365                data,
366                first_seen,
367                last_seen
368            ],
369        )?;
370        tx.execute(
371            "UPDATE relationships SET last_seen = ?4, data = ?5
372             WHERE source_id = ?1 AND target_id = ?2 AND rel_type = ?3",
373            params![
374                edge.source_id,
375                edge.target_id,
376                edge.kind.to_string(),
377                last_seen,
378                data
379            ],
380        )?;
381        Ok(())
382    }
383
384    fn insert_stub_target_for_finding(
385        tx: &Transaction,
386        target_id: &str,
387        finding: &secfinding::Finding,
388    ) -> Result<(), SqliteError> {
389        // Derive NodeType from the target_id's prefix (set by
390        // target_id_from_finding). Default to Endpoint for anything
391        // not in the known set so we don't lose the row.
392        let kind = match target_id.split_once(':').map(|(p, _)| p) {
393            Some("domain") => NodeType::Domain,
394            Some("host") => NodeType::Ip,
395            Some("service") => NodeType::Service,
396            Some("web") => NodeType::Endpoint,
397            _ => NodeType::Endpoint,
398        };
399        let label = finding.target().to_string();
400        let now = std::time::SystemTime::now()
401            .duration_since(std::time::UNIX_EPOCH)
402            .unwrap_or_default()
403            .as_millis() as u64;
404        let now_dt = Self::ms_to_datetime(now);
405        // INSERT OR IGNORE only — never clobber a real Target row that
406        // was upserted earlier in the same transaction with a real
407        // payload.
408        tx.execute(
409            "INSERT OR IGNORE INTO targets (id, kind, label, data, first_seen, last_seen)
410             VALUES (?1, ?2, ?3, '', ?4, ?5)",
411            params![target_id, kind.to_string(), label, now_dt, now_dt],
412        )?;
413        Ok(())
414    }
415
416    fn insert_inferred_target_edges(
417        tx: &Transaction,
418        target: &gossan_core::Target,
419    ) -> Result<(), SqliteError> {
420        match target {
421            gossan_core::Target::Host(h) => {
422                if let Some(domain) = &h.domain {
423                    let e = Edge::new(
424                        format!("domain:{domain}"),
425                        format!("host:{}", h.ip),
426                        EdgeType::ResolvesTo,
427                    );
428                    Self::upsert_edge(tx, &e)?;
429                }
430            }
431            gossan_core::Target::Service(s) => {
432                let e = Edge::new(
433                    format!("host:{}", s.host.ip),
434                    format!("service:{}:{}", s.host.ip, s.port),
435                    EdgeType::HasService,
436                );
437                Self::upsert_edge(tx, &e)?;
438                if let Some(domain) = &s.host.domain {
439                    let e = Edge::new(
440                        format!("domain:{domain}"),
441                        format!("service:{}:{}", s.host.ip, s.port),
442                        EdgeType::HasService,
443                    );
444                    Self::upsert_edge(tx, &e)?;
445                }
446            }
447            gossan_core::Target::Web(w) => {
448                let e = Edge::new(
449                    format!("service:{}:{}", w.service.host.ip, w.service.port),
450                    format!("web:{}", w.url),
451                    EdgeType::Exposes,
452                );
453                Self::upsert_edge(tx, &e)?;
454            }
455            _ => {}
456        }
457        Ok(())
458    }
459
460    /// Raw SQL escape hatch for advanced queries.
461    pub fn conn(&self) -> &Connection {
462        &self.conn
463    }
464}
465
466impl GraphBackend for SqliteBackend {
467    type Error = SqliteError;
468
469    fn init(&mut self) -> Result<(), Self::Error> {
470        self.init_schema()?;
471        Ok(())
472    }
473
474    fn write_nodes(&mut self, nodes: &[Node]) -> Result<(), Self::Error> {
475        let tx = self.conn.transaction()?;
476        for n in nodes {
477            Self::upsert_node(&tx, n)?;
478        }
479        tx.commit()?;
480        Ok(())
481    }
482
483    fn write_edges(&mut self, edges: &[Edge]) -> Result<(), Self::Error> {
484        let tx = self.conn.transaction()?;
485        for e in edges {
486            Self::upsert_edge(&tx, e)?;
487        }
488        tx.commit()?;
489        Ok(())
490    }
491
492    fn read_nodes(&self) -> Result<Vec<Node>, Self::Error> {
493        let mut nodes = Vec::new();
494
495        // Read from targets table
496        let mut stmt = self
497            .conn
498            .prepare("SELECT id, kind, label, data, first_seen, last_seen FROM targets")?;
499        let target_rows = stmt.query_map([], |row| {
500            let kind_str: String = row.get(1)?;
501            let data_str: String = row.get(3)?;
502            let first_seen_str: String = row.get(4)?;
503            let last_seen_str: String = row.get(5)?;
504            Ok(Node {
505                id: row.get(0)?,
506                kind: parse_node_type(&kind_str).unwrap_or(NodeType::Domain),
507                label: row.get(2)?,
508                payload: if data_str.is_empty() {
509                    None
510                } else {
511                    serde_json::from_str(&data_str).ok()
512                },
513                first_seen_ms: Self::datetime_to_ms(&first_seen_str),
514                last_seen_ms: Self::datetime_to_ms(&last_seen_str),
515            })
516        })?;
517        for row in target_rows {
518            nodes.push(row?);
519        }
520
521        // Read from findings table
522        let mut stmt = self
523            .conn
524            .prepare("SELECT id, kind, label, data, first_seen, last_seen FROM findings")?;
525        let finding_rows = stmt.query_map([], |row| {
526            let kind_str: String = row.get(1)?;
527            let data_str: String = row.get(3)?;
528            let first_seen_str: String = row.get(4)?;
529            let last_seen_str: String = row.get(5)?;
530            Ok(Node {
531                id: row.get(0)?,
532                kind: parse_node_type(&kind_str).unwrap_or(NodeType::Finding),
533                label: row.get(2)?,
534                payload: if data_str.is_empty() {
535                    None
536                } else {
537                    serde_json::from_str(&data_str).ok()
538                },
539                first_seen_ms: Self::datetime_to_ms(&first_seen_str),
540                last_seen_ms: Self::datetime_to_ms(&last_seen_str),
541            })
542        })?;
543        for row in finding_rows {
544            nodes.push(row?);
545        }
546
547        Ok(nodes)
548    }
549
550    fn read_edges(&self) -> Result<Vec<Edge>, Self::Error> {
551        let mut stmt = self.conn.prepare(
552            "SELECT source_id, target_id, rel_type, data, first_seen, last_seen FROM relationships",
553        )?;
554        let rows = stmt.query_map([], |row| {
555            let kind_str: String = row.get(2)?;
556            let data_str: String = row.get(3)?;
557            let first_seen_str: String = row.get(4)?;
558            let last_seen_str: String = row.get(5)?;
559            Ok(Edge {
560                source_id: row.get(0)?,
561                target_id: row.get(1)?,
562                kind: parse_edge_type(&kind_str).unwrap_or(EdgeType::HasFinding),
563                payload: if data_str.is_empty() {
564                    None
565                } else {
566                    serde_json::from_str(&data_str).ok()
567                },
568                first_seen_ms: Self::datetime_to_ms(&first_seen_str),
569                last_seen_ms: Self::datetime_to_ms(&last_seen_str),
570            })
571        })?;
572        rows.collect::<Result<Vec<_>, _>>().map_err(Into::into)
573    }
574
575    fn find_nodes_by_type(&self, kind: NodeType) -> Result<Vec<Node>, Self::Error> {
576        let kind_str = kind.to_string();
577        let table = if kind == NodeType::Finding {
578            "findings"
579        } else {
580            "targets"
581        };
582
583        let query = format!(
584            "SELECT id, kind, label, data, first_seen, last_seen FROM {} WHERE kind = ?1",
585            table
586        );
587        let mut stmt = self.conn.prepare(&query)?;
588        let rows = stmt.query_map(params![kind_str], |row| {
589            let data_str: String = row.get(3)?;
590            let first_seen_str: String = row.get(4)?;
591            let last_seen_str: String = row.get(5)?;
592            Ok(Node {
593                id: row.get(0)?,
594                kind: kind.clone(),
595                label: row.get(2)?,
596                payload: if data_str.is_empty() {
597                    None
598                } else {
599                    serde_json::from_str(&data_str).ok()
600                },
601                first_seen_ms: Self::datetime_to_ms(&first_seen_str),
602                last_seen_ms: Self::datetime_to_ms(&last_seen_str),
603            })
604        })?;
605        rows.collect::<Result<Vec<_>, _>>().map_err(Into::into)
606    }
607
608    fn neighbors(
609        &self,
610        node_id: &str,
611        edge_type: Option<EdgeType>,
612    ) -> Result<Vec<Edge>, Self::Error> {
613        let node_id = node_id.to_string();
614        let mut stmt = match edge_type {
615            Some(ref et) => self.conn.prepare(
616                "SELECT source_id, target_id, rel_type, data, first_seen, last_seen
617                 FROM relationships WHERE source_id = ?1 AND rel_type = ?2",
618            )?,
619            None => self.conn.prepare(
620                "SELECT source_id, target_id, rel_type, data, first_seen, last_seen
621                 FROM relationships WHERE source_id = ?1",
622            )?,
623        };
624        let map_row = |row: &rusqlite::Row<'_>| -> Result<Edge, rusqlite::Error> {
625            let kind_str: String = row.get(2)?;
626            let data_str: String = row.get(3)?;
627            let first_seen_str: String = row.get(4)?;
628            let last_seen_str: String = row.get(5)?;
629            Ok(Edge {
630                source_id: row.get(0)?,
631                target_id: row.get(1)?,
632                kind: parse_edge_type(&kind_str).unwrap_or(EdgeType::HasFinding),
633                payload: if data_str.is_empty() {
634                    None
635                } else {
636                    serde_json::from_str(&data_str).ok()
637                },
638                first_seen_ms: Self::datetime_to_ms(&first_seen_str),
639                last_seen_ms: Self::datetime_to_ms(&last_seen_str),
640            })
641        };
642        let rows = match edge_type {
643            Some(et) => stmt.query_map(params![node_id, et.to_string()], map_row)?,
644            None => stmt.query_map(params![node_id], map_row)?,
645        };
646        rows.collect::<Result<Vec<_>, _>>().map_err(Into::into)
647    }
648
649    fn clear(&mut self) -> Result<(), Self::Error> {
650        self.conn.execute("DELETE FROM relationships", [])?;
651        self.conn.execute("DELETE FROM findings", [])?;
652        self.conn.execute("DELETE FROM targets", [])?;
653        Ok(())
654    }
655}
656
657/// Error type for SQLite backend operations.
658#[derive(Debug, thiserror::Error)]
659pub enum SqliteError {
660    #[error("SQLite error: {0}")]
661    Sqlite(#[from] rusqlite::Error),
662    #[error("JSON error: {0}")]
663    Json(#[from] serde_json::Error),
664    #[error("Schema error: {0}")]
665    Schema(String),
666}
667
668/// Compute the deterministic node ID for a `Target`. Used internally
669/// by the SQLite backend (as the `target_id` column on `edges` /
670/// `findings`) and by external callers that want to round-trip
671/// Target identity through a graph store. Promoted to `pub` so the
672/// legendary unit test can pin the ID format
673/// (`domain:<host>` / `host:<ip>` / `service:<ip>:<port>` / etc.)
674/// without depending on the `target_id_from_finding` flavour, which
675/// works on findings instead of targets.
676pub fn target_id(target: &gossan_core::Target) -> String {
677    match target {
678        gossan_core::Target::Domain(d) => format!("domain:{}", d.domain),
679        gossan_core::Target::Host(h) => format!("host:{}", h.ip),
680        gossan_core::Target::Service(s) => format!("service:{}:{}", s.host.ip, s.port),
681        gossan_core::Target::Web(w) => format!("web:{}", w.url),
682        gossan_core::Target::Network(n) => format!("network:{}", n.cidr),
683        gossan_core::Target::Repository(r) => format!("repo:{}", r.url),
684        gossan_core::Target::InternalPackage(p) => format!("pkg:{}", p.name),
685        _ => {
686            let data = serde_json::to_string(target).unwrap_or_default();
687            format!("unknown:{}", &data[..data.len().min(120)])
688        }
689    }
690}
691
692fn target_to_node(target: &gossan_core::Target) -> Node {
693    let id = target_id(target);
694    let (kind, label) = match target {
695        gossan_core::Target::Domain(d) => (NodeType::Domain, d.domain.clone()),
696        gossan_core::Target::Host(h) => (NodeType::Ip, h.ip.to_string()),
697        gossan_core::Target::Service(s) => (NodeType::Service, format!("{}:{}", s.host.ip, s.port)),
698        gossan_core::Target::Web(w) => (NodeType::Endpoint, w.url.to_string()),
699        gossan_core::Target::Network(n) => (NodeType::Ip, n.cidr.clone()),
700        gossan_core::Target::Repository(r) => (NodeType::Endpoint, r.url.to_string()),
701        gossan_core::Target::InternalPackage(p) => (NodeType::Endpoint, p.name.clone()),
702        _ => (NodeType::Endpoint, id.clone()),
703    };
704    Node::new(id, kind, label).with_payload(target)
705}
706
707fn finding_id(finding: &secfinding::Finding) -> String {
708    let namespace = uuid::Uuid::NAMESPACE_OID;
709    let content = format!(
710        "{}:{}:{:?}:{}",
711        finding.target(),
712        finding.title(),
713        finding.severity(),
714        finding.detail()
715    );
716    let id = uuid::Uuid::new_v5(&namespace, content.as_bytes());
717    format!("finding:{id}")
718}
719
720fn finding_to_node(finding: &secfinding::Finding) -> Node {
721    let id = finding_id(finding);
722    Node::new(id, NodeType::Finding, finding.title().to_string()).with_payload(finding)
723}
724
725/// Derive a target node id from a finding target string.
726///
727/// # Errors
728///
729/// Returns an error if the target string cannot be parsed into a known shape.
730pub fn target_id_from_finding(finding: &secfinding::Finding) -> Result<String, SqliteError> {
731    let t = finding.target();
732
733    // Try URL first — but only treat it as a Web target if the parser
734    // actually saw a recognized HTTP-family scheme. `url::Url::parse`
735    // is happy to interpret `"example.com:443"` as `scheme=example.com,
736    // path=443`, which would misclassify a bare host:port pair as a
737    // Web URL.
738    if let Ok(url) = url::Url::parse(t) {
739        if matches!(url.scheme(), "http" | "https" | "ws" | "wss" | "ftp") {
740            return Ok(format!("web:{}", url));
741        }
742    }
743
744    // Try IP address (IPv4 and IPv6)
745    if t.parse::<std::net::IpAddr>().is_ok() {
746        return Ok(format!("host:{t}"));
747    }
748
749    // Try bracketed IPv6
750    if t.starts_with('[') && t.contains("]:") {
751        if let Some(idx) = t.find(']') {
752            let ip_part = &t[1..idx];
753            if ip_part.parse::<std::net::IpAddr>().is_ok() {
754                return Ok(format!("service:{t}"));
755            }
756        }
757    }
758
759    // host:port or ip:port — but avoid misclassifying domains like example.com:443
760    if let Some((host, port)) = t.rsplit_once(':') {
761        if port.parse::<u16>().is_ok() {
762            if host.parse::<std::net::IpAddr>().is_ok() {
763                return Ok(format!("service:{t}"));
764            }
765            // If it looks like a bare IPv6 without brackets, reject rather than guess.
766            if host.contains(':') {
767                return Err(SqliteError::Schema(format!(
768                    "ambiguous IPv6 service target without brackets: {t}"
769                )));
770            }
771        }
772    }
773
774    // Default: domain
775    Ok(format!("domain:{t}"))
776}
777
778fn parse_node_type(s: &str) -> Option<NodeType> {
779    match s {
780        "domain" => Some(NodeType::Domain),
781        "subdomain" => Some(NodeType::Subdomain),
782        "ip" => Some(NodeType::Ip),
783        "port" => Some(NodeType::Port),
784        "service" => Some(NodeType::Service),
785        "tech" => Some(NodeType::Tech),
786        "endpoint" => Some(NodeType::Endpoint),
787        "secret" => Some(NodeType::Secret),
788        "cloud" => Some(NodeType::Cloud),
789        "finding" => Some(NodeType::Finding),
790        _ => None,
791    }
792}
793
794fn parse_edge_type(s: &str) -> Option<EdgeType> {
795    match s {
796        "RESOLVES_TO" => Some(EdgeType::ResolvesTo),
797        "HOSTS" => Some(EdgeType::Hosts),
798        "RUNS" => Some(EdgeType::Runs),
799        "EXPOSES" => Some(EdgeType::Exposes),
800        "LEAKS" => Some(EdgeType::Leaks),
801        "MISCONFIGURED" => Some(EdgeType::Misconfigured),
802        "HAS_FINDING" => Some(EdgeType::HasFinding),
803        "HAS_SERVICE" => Some(EdgeType::HasService),
804        _ => None,
805    }
806}
807
808#[cfg(test)]
809mod tests {
810    use super::*;
811
812    #[test]
813    fn target_id_from_finding_url() {
814        let f = secfinding::Finding::new(
815            "s",
816            "https://example.com/path",
817            secfinding::Severity::Info,
818            "t",
819            "",
820        )
821        .unwrap();
822        assert_eq!(
823            target_id_from_finding(&f).unwrap(),
824            "web:https://example.com/path"
825        );
826    }
827
828    #[test]
829    fn target_id_from_finding_ipv4() {
830        let f =
831            secfinding::Finding::new("s", "1.2.3.4", secfinding::Severity::Info, "t", "").unwrap();
832        assert_eq!(target_id_from_finding(&f).unwrap(), "host:1.2.3.4");
833    }
834
835    #[test]
836    fn target_id_from_finding_ipv6() {
837        let f = secfinding::Finding::new("s", "::1", secfinding::Severity::Info, "t", "").unwrap();
838        assert_eq!(target_id_from_finding(&f).unwrap(), "host:::1");
839    }
840
841    #[test]
842    fn target_id_from_finding_service() {
843        let f = secfinding::Finding::new("s", "1.2.3.4:443", secfinding::Severity::Info, "t", "")
844            .unwrap();
845        assert_eq!(target_id_from_finding(&f).unwrap(), "service:1.2.3.4:443");
846    }
847
848    #[test]
849    fn target_id_from_finding_domain_with_port() {
850        let f =
851            secfinding::Finding::new("s", "example.com:443", secfinding::Severity::Info, "t", "")
852                .unwrap();
853        // Domain with port but no scheme falls through to domain.
854        assert_eq!(
855            target_id_from_finding(&f).unwrap(),
856            "domain:example.com:443"
857        );
858    }
859
860    #[test]
861    fn sqlite_roundtrip() {
862        let mut backend = SqliteBackend::open_in_memory().unwrap();
863        let node = Node::new("n1", NodeType::Domain, "example.com");
864        backend.write_nodes(&[node]).unwrap();
865
866        let edge = Edge::new("n1", "n2", EdgeType::ResolvesTo);
867        backend.write_edges(&[edge]).unwrap();
868
869        let nodes = backend.read_nodes().unwrap();
870        assert_eq!(nodes.len(), 1);
871        assert_eq!(nodes[0].id, "n1");
872
873        let edges = backend.read_edges().unwrap();
874        assert_eq!(edges.len(), 1);
875        assert_eq!(edges[0].kind, EdgeType::ResolvesTo);
876    }
877
878    #[test]
879    fn schema_version_tracked() {
880        let backend = SqliteBackend::open_in_memory().unwrap();
881        let v: i64 = backend
882            .conn()
883            .query_row(
884                "SELECT version FROM schema_version ORDER BY version DESC LIMIT 1",
885                [],
886                |row| row.get(0),
887            )
888            .unwrap();
889        assert_eq!(v, i64::from(SCHEMA_VERSION));
890    }
891}