Skip to main content

stygian_graph/adapters/
storage.rs

1//! Storage adapters — persist and retrieve pipeline [`StorageRecord`](crate::ports::storage::StorageRecord)s.
2//!
3//! # Adapters
4//!
5//! | Adapter | Availability | Backing store |
6//! | --------- | -------------- | --------------- |
7//! | [`NullStorage`](storage::NullStorage) | always | no-op (tests / dry-run) |
8//! | [`FileStorage`](storage::FileStorage) | always | `.jsonl` files on local disk |
9//! | `PostgresStorage` (`storage::postgres::PostgresStorage`) | `feature = "postgres"` | PostgreSQL via sqlx |
10
11use crate::domain::error::{Result, ServiceError, StygianError};
12use crate::ports::storage::{StoragePort, StorageRecord};
13use async_trait::async_trait;
14use std::path::PathBuf;
15use tokio::io::AsyncWriteExt;
16
17// ─────────────────────────────────────────────────────────────────────────────
18// NullStorage
19// ─────────────────────────────────────────────────────────────────────────────
20
21/// No-op storage adapter — discards all records.
22///
23/// Useful for dry-run mode and unit tests where persistence is not required.
24///
25/// # Example
26///
27/// ```
28/// use stygian_graph::adapters::storage::NullStorage;
29/// use stygian_graph::ports::storage::{StoragePort, StorageRecord};
30/// use serde_json::json;
31///
32/// # tokio_test::block_on(async {
33/// let s = NullStorage;
34/// s.store(StorageRecord::new("p", "n", json!(null))).await.unwrap();
35/// let result = s.retrieve("any-id").await.unwrap();
36/// assert!(result.is_none());
37/// # });
38/// ```
39pub struct NullStorage;
40
41#[async_trait]
42impl StoragePort for NullStorage {
43    async fn store(&self, _record: StorageRecord) -> Result<()> {
44        Ok(())
45    }
46
47    async fn retrieve(&self, _id: &str) -> Result<Option<StorageRecord>> {
48        Ok(None)
49    }
50
51    async fn list(&self, _pipeline_id: &str) -> Result<Vec<StorageRecord>> {
52        Ok(vec![])
53    }
54
55    async fn delete(&self, _id: &str) -> Result<()> {
56        Ok(())
57    }
58}
59
60// ─────────────────────────────────────────────────────────────────────────────
61// FileStorage
62// ─────────────────────────────────────────────────────────────────────────────
63
64/// File-based storage adapter — one `.jsonl` file per pipeline.
65///
66/// Each pipeline gets its own file: `<dir>/<pipeline_id>.jsonl`.
67/// Records are appended one JSON object per line.
68///
69/// # Example
70///
71/// ```no_run
72/// use stygian_graph::adapters::storage::FileStorage;
73/// use stygian_graph::ports::storage::{StoragePort, StorageRecord};
74/// use serde_json::json;
75/// use std::path::PathBuf;
76///
77/// # tokio_test::block_on(async {
78/// let storage = FileStorage::new(PathBuf::from("/tmp/stygian-results"));
79/// let r = StorageRecord::new("pipe-1", "fetch", json!({"url": "https://example.com"}));
80/// storage.store(r).await.unwrap();
81/// # });
82/// ```
83pub struct FileStorage {
84    dir: PathBuf,
85}
86
87impl FileStorage {
88    /// Create a [`FileStorage`] backed by `dir`.
89    ///
90    /// The directory will be created on first write if it does not exist.
91    ///
92    /// # Example
93    ///
94    /// ```
95    /// use stygian_graph::adapters::storage::FileStorage;
96    /// use std::path::PathBuf;
97    ///
98    /// let s = FileStorage::new(PathBuf::from("/tmp/data"));
99    /// ```
100    #[must_use]
101    pub const fn new(dir: PathBuf) -> Self {
102        Self { dir }
103    }
104
105    fn pipeline_file(&self, pipeline_id: &str) -> PathBuf {
106        // Sanitise: replace path separators so callers cannot escape the dir
107        let safe_id = pipeline_id.replace(['/', '\\', '.', ':'], "_");
108        self.dir.join(format!("{safe_id}.jsonl"))
109    }
110}
111
112#[async_trait]
113impl StoragePort for FileStorage {
114    async fn store(&self, record: StorageRecord) -> Result<()> {
115        tokio::fs::create_dir_all(&self.dir).await.map_err(|e| {
116            StygianError::Service(ServiceError::InvalidResponse(format!(
117                "FileStorage: create_dir_all failed: {e}"
118            )))
119        })?;
120
121        let path = self.pipeline_file(&record.pipeline_id);
122        let mut line = serde_json::to_string(&record).map_err(|e| {
123            StygianError::Service(ServiceError::InvalidResponse(format!(
124                "FileStorage: serialise record failed: {e}"
125            )))
126        })?;
127        line.push('\n');
128
129        let mut file = tokio::fs::OpenOptions::new()
130            .create(true)
131            .append(true)
132            .open(&path)
133            .await
134            .map_err(|e| {
135                StygianError::Service(ServiceError::InvalidResponse(format!(
136                    "FileStorage: open {}: {e}",
137                    path.display()
138                )))
139            })?;
140
141        file.write_all(line.as_bytes()).await.map_err(|e| {
142            StygianError::Service(ServiceError::InvalidResponse(format!(
143                "FileStorage: write failed: {e}"
144            )))
145        })?;
146
147        // Flush + sync before dropping the file handle so a subsequent
148        // `list()` on the same path in another task observes the
149        // appended line. Without this, tokio's async File drop can race
150        // the reader's `read_to_string` on filesystems that buffer
151        // append-mode writes (notably CI runners with networked/overlay
152        // filesystems).
153        file.flush().await.map_err(|e| {
154            StygianError::Service(ServiceError::InvalidResponse(format!(
155                "FileStorage: flush failed: {e}"
156            )))
157        })?;
158        file.sync_all().await.map_err(|e| {
159            StygianError::Service(ServiceError::InvalidResponse(format!(
160                "FileStorage: sync_all failed: {e}"
161            )))
162        })?;
163
164        Ok(())
165    }
166
167    async fn retrieve(&self, id: &str) -> Result<Option<StorageRecord>> {
168        // Scan all .jsonl files — linear scan is acceptable for moderate volumes
169        let Ok(mut dir) = tokio::fs::read_dir(&self.dir).await else {
170            return Ok(None);
171        };
172
173        while let Ok(Some(entry)) = dir.next_entry().await {
174            let path = entry.path();
175            if path.extension().and_then(|e| e.to_str()) != Some("jsonl") {
176                continue;
177            }
178            let Ok(content) = tokio::fs::read_to_string(&path).await else {
179                continue;
180            };
181            for line in content.lines() {
182                if let Ok(record) = serde_json::from_str::<StorageRecord>(line)
183                    && record.id == id
184                {
185                    return Ok(Some(record));
186                }
187            }
188        }
189
190        Ok(None)
191    }
192
193    async fn list(&self, pipeline_id: &str) -> Result<Vec<StorageRecord>> {
194        let path = self.pipeline_file(pipeline_id);
195        let Ok(content) = tokio::fs::read_to_string(&path).await else {
196            return Ok(vec![]);
197        };
198
199        let records = content
200            .lines()
201            .filter(|l| !l.is_empty())
202            .filter_map(|line| serde_json::from_str::<StorageRecord>(line).ok())
203            .collect();
204
205        Ok(records)
206    }
207
208    async fn delete(&self, id: &str) -> Result<()> {
209        // Read-filter-rewrite strategy (adequate for typical pipeline sizes)
210        let Ok(mut dir) = tokio::fs::read_dir(&self.dir).await else {
211            return Ok(()); // dir does not exist → nothing to delete
212        };
213
214        while let Ok(Some(entry)) = dir.next_entry().await {
215            let path = entry.path();
216            if path.extension().and_then(|e| e.to_str()) != Some("jsonl") {
217                continue;
218            }
219            let Ok(content) = tokio::fs::read_to_string(&path).await else {
220                continue;
221            };
222
223            let (kept, found): (Vec<&str>, bool) = {
224                let mut found = false;
225                let kept = content
226                    .lines()
227                    .filter(|line| {
228                        if let Ok(r) = serde_json::from_str::<StorageRecord>(line)
229                            && r.id == id
230                        {
231                            found = true;
232                            return false;
233                        }
234                        true
235                    })
236                    .collect::<Vec<_>>();
237                (kept, found)
238            };
239
240            if found {
241                let new_content = kept.join("\n");
242                let new_content = if new_content.is_empty() {
243                    new_content
244                } else {
245                    format!("{new_content}\n")
246                };
247                tokio::fs::write(&path, new_content.as_bytes())
248                    .await
249                    .map_err(|e| {
250                        StygianError::Service(ServiceError::InvalidResponse(format!(
251                            "FileStorage: rewrite after delete failed: {e}"
252                        )))
253                    })?;
254                return Ok(());
255            }
256        }
257
258        Ok(())
259    }
260}
261
262// ─────────────────────────────────────────────────────────────────────────────
263// PostgresStorage — feature = "postgres"
264// ─────────────────────────────────────────────────────────────────────────────
265
266#[cfg(feature = "postgres")]
267pub use postgres::PostgresStorage;
268
269#[cfg(feature = "postgres")]
270mod postgres {
271    //! PostgreSQL-backed storage via sqlx.
272    //!
273    //! Assumes the following table exists:
274    //!
275    //! ```sql
276    //! CREATE TABLE IF NOT EXISTS pipeline_records (
277    //!     id          TEXT PRIMARY KEY,
278    //!     pipeline_id TEXT NOT NULL,
279    //!     node_name   TEXT NOT NULL,
280    //!     data        JSONB NOT NULL,
281    //!     metadata    JSONB NOT NULL DEFAULT '{}',
282    //!     timestamp_ms BIGINT NOT NULL
283    //! );
284    //! CREATE INDEX IF NOT EXISTS idx_pipeline_records_pipeline_id
285    //!     ON pipeline_records (pipeline_id);
286    //! ```
287
288    use crate::domain::error::{Result, ServiceError, StygianError};
289    use crate::ports::storage::{StoragePort, StorageRecord};
290    use sqlx::{PgPool, Row};
291
292    /// `PostgreSQL` storage adapter.
293    ///
294    /// # Example
295    ///
296    /// ```no_run
297    /// use stygian_graph::adapters::storage::PostgresStorage;
298    /// use sqlx::PgPool;
299    ///
300    /// # tokio_test::block_on(async {
301    /// let pool = PgPool::connect("postgres://localhost/stygian").await.unwrap();
302    /// let storage = PostgresStorage::new(pool);
303    /// # });
304    /// ```
305    pub struct PostgresStorage {
306        pool: PgPool,
307    }
308
309    impl PostgresStorage {
310        /// Create a new [`PostgresStorage`] from a connection pool.
311        ///
312        /// # Example
313        ///
314        /// ```no_run
315        /// use stygian_graph::adapters::storage::PostgresStorage;
316        /// use sqlx::PgPool;
317        ///
318        /// # tokio_test::block_on(async {
319        /// let pool = PgPool::connect("postgres://localhost/stygian").await.unwrap();
320        /// let s = PostgresStorage::new(pool);
321        /// # });
322        /// ```
323        #[must_use]
324        pub const fn new(pool: PgPool) -> Self {
325            Self { pool }
326        }
327    }
328
329    #[async_trait::async_trait]
330    impl StoragePort for PostgresStorage {
331        async fn store(&self, record: StorageRecord) -> Result<()> {
332            let metadata_json = serde_json::to_value(&record.metadata).map_err(|e| {
333                StygianError::Service(ServiceError::InvalidResponse(format!(
334                    "PostgresStorage: metadata serialise: {e}"
335                )))
336            })?;
337
338            sqlx::query(
339                "
340                INSERT INTO pipeline_records
341                    (id, pipeline_id, node_name, data, metadata, timestamp_ms)
342                VALUES ($1, $2, $3, $4, $5, $6)
343                ON CONFLICT (id) DO NOTHING
344                ",
345            )
346            .bind(&record.id)
347            .bind(&record.pipeline_id)
348            .bind(&record.node_name)
349            .bind(&record.data)
350            .bind(metadata_json)
351            .bind(i64::try_from(record.timestamp_ms).unwrap_or(i64::MAX))
352            .execute(&self.pool)
353            .await
354            .map_err(|e| {
355                StygianError::Service(ServiceError::InvalidResponse(format!(
356                    "PostgresStorage: insert failed: {e}"
357                )))
358            })?;
359
360            Ok(())
361        }
362
363        async fn retrieve(&self, id: &str) -> Result<Option<StorageRecord>> {
364            let row = sqlx::query(
365                "
366                SELECT id, pipeline_id, node_name, data, metadata, timestamp_ms
367                FROM pipeline_records
368                WHERE id = $1
369                ",
370            )
371            .bind(id)
372            .fetch_optional(&self.pool)
373            .await
374            .map_err(|e| {
375                StygianError::Service(ServiceError::InvalidResponse(format!(
376                    "PostgresStorage: retrieve failed: {e}"
377                )))
378            })?;
379
380            row.map_or(Ok(None), |r| {
381                let metadata = serde_json::from_value(r.get::<serde_json::Value, _>("metadata"))
382                    .unwrap_or_default();
383                Ok(Some(StorageRecord {
384                    id: r.get("id"),
385                    pipeline_id: r.get("pipeline_id"),
386                    node_name: r.get("node_name"),
387                    data: r.get("data"),
388                    metadata,
389                    timestamp_ms: u64::try_from(r.get::<i64, _>("timestamp_ms")).unwrap_or(0),
390                }))
391            })
392        }
393
394        async fn list(&self, pipeline_id: &str) -> Result<Vec<StorageRecord>> {
395            let rows = sqlx::query(
396                "
397                SELECT id, pipeline_id, node_name, data, metadata, timestamp_ms
398                FROM pipeline_records
399                WHERE pipeline_id = $1
400                ORDER BY timestamp_ms ASC
401                ",
402            )
403            .bind(pipeline_id)
404            .fetch_all(&self.pool)
405            .await
406            .map_err(|e| {
407                StygianError::Service(ServiceError::InvalidResponse(format!(
408                    "PostgresStorage: list failed: {e}"
409                )))
410            })?;
411
412            let records = rows
413                .into_iter()
414                .map(|r| {
415                    let metadata =
416                        serde_json::from_value(r.get::<serde_json::Value, _>("metadata"))
417                            .unwrap_or_default();
418                    StorageRecord {
419                        id: r.get("id"),
420                        pipeline_id: r.get("pipeline_id"),
421                        node_name: r.get("node_name"),
422                        data: r.get("data"),
423                        metadata,
424                        timestamp_ms: u64::try_from(r.get::<i64, _>("timestamp_ms")).unwrap_or(0),
425                    }
426                })
427                .collect();
428
429            Ok(records)
430        }
431
432        async fn delete(&self, id: &str) -> Result<()> {
433            sqlx::query("DELETE FROM pipeline_records WHERE id = $1")
434                .bind(id)
435                .execute(&self.pool)
436                .await
437                .map_err(|e| {
438                    StygianError::Service(ServiceError::InvalidResponse(format!(
439                        "PostgresStorage: delete failed: {e}"
440                    )))
441                })?;
442
443            Ok(())
444        }
445    }
446}
447
448// ─────────────────────────────────────────────────────────────────────────────
449// Tests
450// ─────────────────────────────────────────────────────────────────────────────
451
452#[cfg(test)]
453#[allow(clippy::unwrap_used, clippy::indexing_slicing)]
454mod tests {
455    use super::{FileStorage, NullStorage};
456    use crate::ports::storage::{StoragePort, StorageRecord};
457    use serde_json::json;
458
459    #[tokio::test]
460    async fn null_storage_store_and_retrieve() {
461        let s = NullStorage;
462        let r = StorageRecord::new("p", "n", json!(null));
463        s.store(r.clone()).await.unwrap();
464        let got = s.retrieve(&r.id).await.unwrap();
465        assert!(got.is_none(), "NullStorage must always return None");
466    }
467
468    #[tokio::test]
469    async fn null_storage_list_and_delete_are_noops() {
470        let s = NullStorage;
471        let list = s.list("any").await.unwrap();
472        assert!(list.is_empty());
473        s.delete("any-id").await.unwrap();
474    }
475
476    #[tokio::test]
477    async fn file_storage_roundtrip() {
478        let dir = tempfile::tempdir().unwrap();
479        let storage = FileStorage::new(dir.path().to_path_buf());
480
481        let r = StorageRecord::new(
482            "pipe-roundtrip",
483            "fetch",
484            json!({"url": "https://example.com"}),
485        );
486        let id = r.id.clone();
487
488        storage.store(r).await.unwrap();
489
490        let retrieved = storage.retrieve(&id).await.unwrap().unwrap();
491        assert_eq!(retrieved.id, id);
492        assert_eq!(retrieved.pipeline_id, "pipe-roundtrip");
493        assert_eq!(retrieved.node_name, "fetch");
494    }
495
496    #[tokio::test]
497    async fn file_storage_list_scoped_to_pipeline() {
498        let dir = tempfile::tempdir().unwrap();
499        let storage = FileStorage::new(dir.path().to_path_buf());
500
501        storage
502            .store(StorageRecord::new("pipe-a", "step1", json!(1)))
503            .await
504            .unwrap();
505        storage
506            .store(StorageRecord::new("pipe-a", "step2", json!(2)))
507            .await
508            .unwrap();
509        storage
510            .store(StorageRecord::new("pipe-b", "step1", json!(3)))
511            .await
512            .unwrap();
513
514        let pipe_a = storage.list("pipe-a").await.unwrap();
515        assert_eq!(pipe_a.len(), 2);
516
517        let pipe_b = storage.list("pipe-b").await.unwrap();
518        assert_eq!(pipe_b.len(), 1);
519    }
520
521    #[tokio::test]
522    async fn file_storage_delete_removes_record() {
523        let dir = tempfile::tempdir().unwrap();
524        let storage = FileStorage::new(dir.path().to_path_buf());
525
526        let r1 = StorageRecord::new("pipe-del", "n", json!(1));
527        let r2 = StorageRecord::new("pipe-del", "n", json!(2));
528        let id1 = r1.id.clone();
529
530        storage.store(r1).await.unwrap();
531        storage.store(r2).await.unwrap();
532
533        storage.delete(&id1).await.unwrap();
534
535        let records = storage.list("pipe-del").await.unwrap();
536        assert_eq!(records.len(), 1);
537        assert_ne!(records[0].id, id1);
538    }
539
540    #[tokio::test]
541    async fn file_storage_retrieve_not_found_returns_none() {
542        let dir = tempfile::tempdir().unwrap();
543        let storage = FileStorage::new(dir.path().to_path_buf());
544        let result = storage.retrieve("no-such-id").await.unwrap();
545        assert!(result.is_none());
546    }
547
548    #[tokio::test]
549    async fn file_storage_path_sanitises_separators() {
550        let dir = tempfile::tempdir().unwrap();
551        let storage = FileStorage::new(dir.path().to_path_buf());
552
553        // pipeline_id with slashes should not escape the base directory
554        let r = StorageRecord::new("../../etc/passwd", "n", json!(null));
555        storage.store(r).await.unwrap();
556
557        let files: Vec<_> = std::fs::read_dir(dir.path())
558            .unwrap()
559            .filter_map(Result::ok)
560            .collect();
561        // File must be inside the temp dir, not some other directory
562        assert_eq!(files.len(), 1);
563        let fname = files[0].file_name();
564        assert!(
565            fname.to_string_lossy().contains("__"),
566            "separators must be sanitised: got {fname:?}"
567        );
568    }
569
570    #[tokio::test]
571    async fn file_storage_retrieve_finds_correct_record() {
572        let dir = tempfile::tempdir().unwrap();
573        let storage = FileStorage::new(dir.path().to_path_buf());
574
575        // Store records across two pipelines to exercise full-dir scan in retrieve
576        let r1 = StorageRecord::new("pipe-x", "node-1", json!({"val": 1}));
577        let r2 = StorageRecord::new("pipe-y", "node-2", json!({"val": 2}));
578        let id1 = r1.id.clone();
579        let id2 = r2.id.clone();
580
581        storage.store(r1).await.unwrap();
582        storage.store(r2).await.unwrap();
583
584        let found = storage.retrieve(&id1).await.unwrap().unwrap();
585        assert_eq!(found.id, id1);
586        assert_eq!(found.pipeline_id, "pipe-x");
587
588        let found2 = storage.retrieve(&id2).await.unwrap().unwrap();
589        assert_eq!(found2.id, id2);
590        assert_eq!(found2.pipeline_id, "pipe-y");
591    }
592
593    #[tokio::test]
594    async fn file_storage_retrieve_missing_returns_none() {
595        let dir = tempfile::tempdir().unwrap();
596        let storage = FileStorage::new(dir.path().to_path_buf());
597        // Store something so the dir exists and the scan loop runs
598        storage
599            .store(StorageRecord::new("p", "n", json!(0)))
600            .await
601            .unwrap();
602        let result = storage.retrieve("nonexistent-id").await.unwrap();
603        assert!(result.is_none());
604    }
605
606    #[tokio::test]
607    async fn file_storage_delete_nonexistent_dir_is_noop() {
608        // Dir is never created — delete should return Ok without panicking
609        let storage = FileStorage::new(std::path::PathBuf::from("/tmp/stygian-no-such-dir-xyz"));
610        storage.delete("any-id").await.unwrap();
611    }
612
613    #[tokio::test]
614    async fn file_storage_delete_id_not_present_is_noop() {
615        let dir = tempfile::tempdir().unwrap();
616        let storage = FileStorage::new(dir.path().to_path_buf());
617        let r = StorageRecord::new("pipe-z", "n", json!(42));
618        storage.store(r).await.unwrap();
619        // Deleting a non-existent id should not modify the file
620        storage.delete("totally-unknown-id").await.unwrap();
621        let records = storage.list("pipe-z").await.unwrap();
622        assert_eq!(records.len(), 1);
623    }
624
625    #[tokio::test]
626    async fn file_storage_list_missing_pipeline_returns_empty() {
627        let dir = tempfile::tempdir().unwrap();
628        let storage = FileStorage::new(dir.path().to_path_buf());
629        let records = storage.list("never-stored").await.unwrap();
630        assert!(records.is_empty());
631    }
632}