millipede_storage_memory/
dataset.rs1use 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
15pub struct MemoryDataset {
17 name: String,
18 inner: Mutex<DatasetState>,
19}
20
21impl MemoryDataset {
22 #[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 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 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 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}