Skip to main content

millipede_storage_memory/
dataset.rs

1use futures_util::{stream, stream::BoxStream};
2use millipede_core::storage::{
3    Dataset, DatasetInfo, ListOptions, Page, StorageError, StorageResult,
4};
5use serde_json::Value;
6use std::{collections::BTreeSet, path::Path, sync::Mutex};
7use time::OffsetDateTime;
8
9struct DatasetState {
10    items: Vec<Value>,
11    created_at: OffsetDateTime,
12    modified_at: OffsetDateTime,
13}
14
15/// An in-process, append-only JSON dataset.
16pub struct MemoryDataset {
17    name: String,
18    inner: Mutex<DatasetState>,
19}
20
21impl MemoryDataset {
22    /// Creates an empty dataset with the supplied name.
23    #[must_use]
24    pub fn new(name: impl Into<String>) -> Self {
25        let now = OffsetDateTime::now_utc();
26        Self {
27            name: name.into(),
28            inner: Mutex::new(DatasetState {
29                items: Vec::new(),
30                created_at: now,
31                modified_at: now,
32            }),
33        }
34    }
35
36    fn lock(&self) -> std::sync::MutexGuard<'_, DatasetState> {
37        // A panic while holding this lock is a programming bug, so poisoning is unrecoverable.
38        self.inner.lock().expect("MemoryDataset mutex poisoned")
39    }
40
41    fn sliced_items(&self, opts: &ListOptions) -> (Vec<Value>, u64) {
42        let state = self.lock();
43        let total = state.items.len() as u64;
44        let items: Vec<_> = if opts.desc {
45            state.items.iter().rev().cloned().collect()
46        } else {
47            state.items.clone()
48        };
49        let items = items
50            .into_iter()
51            .skip(usize::try_from(opts.offset).unwrap_or(usize::MAX))
52            .take(opts.limit.map_or(usize::MAX, |limit| {
53                usize::try_from(limit).unwrap_or(usize::MAX)
54            }))
55            .collect();
56        (items, total)
57    }
58
59    pub(crate) fn clear(&self) {
60        let mut state = self.lock();
61        state.items.clear();
62        state.modified_at = OffsetDateTime::now_utc();
63    }
64}
65
66#[async_trait::async_trait]
67impl Dataset for MemoryDataset {
68    async fn push_json(&self, item: Value) -> StorageResult<()> {
69        let mut state = self.lock();
70        state.items.push(item);
71        state.modified_at = OffsetDateTime::now_utc();
72        Ok(())
73    }
74
75    async fn push_json_batch(&self, items: Vec<Value>) -> StorageResult<()> {
76        let mut state = self.lock();
77        state.items.extend(items);
78        state.modified_at = OffsetDateTime::now_utc();
79        Ok(())
80    }
81
82    async fn list_raw(&self, opts: ListOptions) -> StorageResult<Page<Value>> {
83        let (items, total) = self.sliced_items(&opts);
84        Ok(Page {
85            items,
86            total,
87            offset: opts.offset,
88            limit: opts.limit,
89        })
90    }
91
92    fn stream_raw(&self, opts: ListOptions) -> BoxStream<'_, StorageResult<Value>> {
93        let (items, _) = self.sliced_items(&opts);
94        Box::pin(stream::iter(items.into_iter().map(Ok)))
95    }
96
97    async fn export_json(&self, path: &Path) -> StorageResult<()> {
98        let items = self.lock().items.clone();
99        let bytes = serde_json::to_vec_pretty(&items)?;
100        // This synchronous write is bounded to an already-materialized in-memory buffer.
101        std::fs::write(path, bytes).map_err(StorageError::Io)
102    }
103
104    async fn export_csv(&self, path: &Path) -> StorageResult<()> {
105        let items = self.lock().items.clone();
106        let mut columns = BTreeSet::new();
107        for item in &items {
108            let object = item.as_object().ok_or(StorageError::Unsupported(
109                "export_csv requires object items",
110            ))?;
111            columns.extend(object.keys().cloned());
112        }
113        let columns: Vec<_> = columns.into_iter().collect();
114        let mut rows = vec![
115            columns
116                .iter()
117                .map(|key| csv_field(key))
118                .collect::<Vec<_>>()
119                .join(","),
120        ];
121        for item in &items {
122            let object = item.as_object().expect("objects validated above");
123            rows.push(
124                columns
125                    .iter()
126                    .map(|key| match object.get(key) {
127                        None => String::new(),
128                        Some(Value::String(value)) => csv_field(value),
129                        Some(value) => csv_field(&value.to_string()),
130                    })
131                    .collect::<Vec<_>>()
132                    .join(","),
133            );
134        }
135        // This synchronous write is bounded to an already-materialized in-memory buffer.
136        std::fs::write(path, rows.join("\r\n")).map_err(StorageError::Io)
137    }
138
139    async fn info(&self) -> StorageResult<DatasetInfo> {
140        let state = self.lock();
141        Ok(DatasetInfo::new(
142            self.name.clone(),
143            state.items.len() as u64,
144            state.created_at,
145            state.modified_at,
146        ))
147    }
148}
149
150fn csv_field(value: &str) -> String {
151    if value.contains([',', '"', '\r', '\n']) {
152        format!("\"{}\"", value.replace('"', "\"\""))
153    } else {
154        value.to_owned()
155    }
156}