1use 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
17pub 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
60pub struct FileStorage {
84 dir: PathBuf,
85}
86
87impl FileStorage {
88 #[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 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 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 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 let Ok(mut dir) = tokio::fs::read_dir(&self.dir).await else {
211 return Ok(()); };
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#[cfg(feature = "postgres")]
267pub use postgres::PostgresStorage;
268
269#[cfg(feature = "postgres")]
270mod postgres {
271 use crate::domain::error::{Result, ServiceError, StygianError};
289 use crate::ports::storage::{StoragePort, StorageRecord};
290 use sqlx::{PgPool, Row};
291
292 pub struct PostgresStorage {
306 pool: PgPool,
307 }
308
309 impl PostgresStorage {
310 #[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#[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 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 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 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 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 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 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}