Skip to main content

backtest_server/
artifact_store.rs

1//! Filesystem-backed storage for large backtest result JSON artifacts.
2
3use std::fs::{self, OpenOptions};
4use std::io::{Read, Seek, SeekFrom, Write};
5use std::path::{Path, PathBuf};
6use std::sync::Mutex;
7use std::time::{Duration, SystemTime};
8
9use nanoid::nanoid;
10use sha2::{Digest, Sha256};
11use thiserror::Error;
12
13use crate::rpc_types::{RESULT_FORMAT_VERSION, ResultArtifactRefMsg};
14
15const ARTIFACT_EXTENSION: &str = "json";
16const PART_EXTENSION: &str = "part";
17const TRANSPORT_LIMIT_BYTES: usize = 16 * 1024 * 1024;
18const TRANSPORT_RESPONSE_RESERVE_BYTES: usize = 64 * 1024;
19
20/// One raw artifact chunk read from disk.
21#[derive(Debug, Clone, PartialEq, Eq)]
22pub struct ArtifactChunk {
23    pub offset: u64,
24    pub bytes: Vec<u8>,
25    pub eof: bool,
26}
27
28/// Errors returned by the artifact store.
29#[derive(Debug, Error)]
30pub enum ArtifactStoreError {
31    #[error("invalid artifact id")]
32    InvalidId,
33    #[error("artifact '{0}' was not found")]
34    NotFound(String),
35    #[error("artifact path escaped the configured store directory")]
36    PathEscape,
37    #[error("artifact offset {offset} exceeds byte length {byte_len}")]
38    InvalidOffset { offset: u64, byte_len: u64 },
39    #[error("artifact is {byte_len} bytes but the store capacity is {max_total_bytes} bytes")]
40    CapacityExceeded { byte_len: u64, max_total_bytes: u64 },
41    #[error(
42        "artifact chunk size {chunk_size} cannot fit in the 16 MiB transport after base64 encoding"
43    )]
44    ChunkTooLarge { chunk_size: usize },
45    #[error("artifact store I/O error: {0}")]
46    Io(#[from] std::io::Error),
47}
48
49/// Server-owned filesystem artifact store.
50#[derive(Debug)]
51pub struct ArtifactStore {
52    root: PathBuf,
53    inline_limit_bytes: usize,
54    chunk_size: usize,
55    retention: Duration,
56    max_total_bytes: u64,
57    lock: Mutex<()>,
58}
59
60impl ArtifactStore {
61    pub fn new(
62        directory: impl AsRef<Path>,
63        inline_limit_bytes: usize,
64        chunk_size: usize,
65        retention: Duration,
66        max_total_bytes: u64,
67    ) -> Result<Self, ArtifactStoreError> {
68        let encoded_chunk_len = chunk_size.div_ceil(3).saturating_mul(4);
69        if chunk_size == 0
70            || encoded_chunk_len
71                > TRANSPORT_LIMIT_BYTES.saturating_sub(TRANSPORT_RESPONSE_RESERVE_BYTES)
72        {
73            return Err(ArtifactStoreError::ChunkTooLarge { chunk_size });
74        }
75
76        fs::create_dir_all(directory.as_ref())?;
77        let root = fs::canonicalize(directory.as_ref())?;
78        let store = Self {
79            root,
80            inline_limit_bytes,
81            chunk_size,
82            retention,
83            max_total_bytes,
84            lock: Mutex::new(()),
85        };
86        store.cleanup_expired()?;
87        Ok(store)
88    }
89
90    pub fn inline_limit_bytes(&self) -> usize {
91        self.inline_limit_bytes
92    }
93
94    pub fn chunk_size(&self) -> usize {
95        self.chunk_size
96    }
97
98    pub fn persist_json(&self, bytes: &[u8]) -> Result<ResultArtifactRefMsg, ArtifactStoreError> {
99        let byte_len = bytes.len() as u64;
100        if byte_len > self.max_total_bytes {
101            return Err(ArtifactStoreError::CapacityExceeded {
102                byte_len,
103                max_total_bytes: self.max_total_bytes,
104            });
105        }
106
107        let _guard = self.lock.lock().unwrap_or_else(|error| error.into_inner());
108        self.cleanup_expired_locked(SystemTime::now())?;
109        self.reserve_capacity_locked(byte_len)?;
110
111        let (artifact_id, part_path, final_path, mut file) = loop {
112            let artifact_id = format!("result_{}", nanoid!(24));
113            let part_path = self.path_for_id_with_extension(&artifact_id, PART_EXTENSION)?;
114            let final_path = self.path_for_id_with_extension(&artifact_id, ARTIFACT_EXTENSION)?;
115            match OpenOptions::new()
116                .write(true)
117                .create_new(true)
118                .open(&part_path)
119            {
120                Ok(file) => break (artifact_id, part_path, final_path, file),
121                Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => continue,
122                Err(error) => return Err(error.into()),
123            }
124        };
125
126        let write_result = (|| -> std::io::Result<()> {
127            file.write_all(bytes)?;
128            file.sync_all()?;
129            drop(file);
130            fs::rename(&part_path, &final_path)
131        })();
132        if let Err(error) = write_result {
133            let _ = fs::remove_file(&part_path);
134            return Err(error.into());
135        }
136
137        Ok(ResultArtifactRefMsg {
138            format_version: RESULT_FORMAT_VERSION,
139            artifact_id,
140            byte_len,
141            sha256: sha256_hex(bytes),
142            chunk_size: self.chunk_size as u64,
143        })
144    }
145
146    pub fn read_chunk(
147        &self,
148        artifact_id: &str,
149        offset: u64,
150    ) -> Result<ArtifactChunk, ArtifactStoreError> {
151        let _guard = self.lock.lock().unwrap_or_else(|error| error.into_inner());
152        let path = self.existing_artifact_path(artifact_id)?;
153        let mut file = OpenOptions::new().read(true).write(true).open(path)?;
154        let byte_len = file.metadata()?.len();
155        if offset > byte_len {
156            return Err(ArtifactStoreError::InvalidOffset { offset, byte_len });
157        }
158
159        let read_len = (byte_len - offset).min(self.chunk_size as u64) as usize;
160        let mut bytes = vec![0; read_len];
161        file.seek(SeekFrom::Start(offset))?;
162        if read_len > 0 {
163            file.read_exact(&mut bytes)?;
164        }
165        file.set_modified(SystemTime::now())?;
166        Ok(ArtifactChunk {
167            offset,
168            eof: offset + bytes.len() as u64 >= byte_len,
169            bytes,
170        })
171    }
172
173    pub fn delete(&self, artifact_id: &str) -> Result<bool, ArtifactStoreError> {
174        let _guard = self.lock.lock().unwrap_or_else(|error| error.into_inner());
175        match self.existing_artifact_path(artifact_id) {
176            Ok(path) => {
177                fs::remove_file(path)?;
178                Ok(true)
179            }
180            Err(ArtifactStoreError::NotFound(_)) => Ok(false),
181            Err(error) => Err(error),
182        }
183    }
184
185    pub fn cleanup_expired(&self) -> Result<usize, ArtifactStoreError> {
186        let _guard = self.lock.lock().unwrap_or_else(|error| error.into_inner());
187        self.cleanup_expired_locked(SystemTime::now())
188    }
189
190    fn reserve_capacity_locked(&self, incoming: u64) -> Result<(), ArtifactStoreError> {
191        let total = self.artifact_bytes_locked()?;
192        if total.saturating_add(incoming) > self.max_total_bytes {
193            return Err(ArtifactStoreError::CapacityExceeded {
194                byte_len: incoming,
195                max_total_bytes: self.max_total_bytes,
196            });
197        }
198        Ok(())
199    }
200
201    fn cleanup_expired_locked(&self, now: SystemTime) -> Result<usize, ArtifactStoreError> {
202        let mut removed = 0;
203        for entry in fs::read_dir(&self.root)? {
204            let entry = entry?;
205            let path = entry.path();
206            let metadata = fs::symlink_metadata(&path)?;
207            if !metadata.file_type().is_file() {
208                continue;
209            }
210            let extension = path.extension().and_then(|value| value.to_str());
211            if extension != Some(ARTIFACT_EXTENSION) && extension != Some(PART_EXTENSION) {
212                continue;
213            }
214            let modified = metadata.modified().unwrap_or(SystemTime::UNIX_EPOCH);
215            let age = now.duration_since(modified).unwrap_or_default();
216            if age >= self.retention {
217                fs::remove_file(path)?;
218                removed += 1;
219            }
220        }
221        Ok(removed)
222    }
223
224    fn artifact_bytes_locked(&self) -> Result<u64, ArtifactStoreError> {
225        let mut total = 0_u64;
226        for entry in fs::read_dir(&self.root)? {
227            let entry = entry?;
228            let path = entry.path();
229            if path.extension().and_then(|value| value.to_str()) != Some(ARTIFACT_EXTENSION) {
230                continue;
231            }
232            let metadata = fs::symlink_metadata(&path)?;
233            if metadata.file_type().is_file() {
234                total = total.saturating_add(metadata.len());
235            }
236        }
237        Ok(total)
238    }
239
240    fn existing_artifact_path(&self, artifact_id: &str) -> Result<PathBuf, ArtifactStoreError> {
241        let path = self.path_for_id_with_extension(artifact_id, ARTIFACT_EXTENSION)?;
242        let canonical = match fs::canonicalize(&path) {
243            Ok(path) => path,
244            Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
245                return Err(ArtifactStoreError::NotFound(artifact_id.to_owned()));
246            }
247            Err(error) => return Err(error.into()),
248        };
249        if canonical.parent() != Some(self.root.as_path()) {
250            return Err(ArtifactStoreError::PathEscape);
251        }
252        if !fs::symlink_metadata(&canonical)?.file_type().is_file() {
253            return Err(ArtifactStoreError::NotFound(artifact_id.to_owned()));
254        }
255        Ok(canonical)
256    }
257
258    fn path_for_id_with_extension(
259        &self,
260        artifact_id: &str,
261        extension: &str,
262    ) -> Result<PathBuf, ArtifactStoreError> {
263        validate_artifact_id(artifact_id)?;
264        Ok(self.root.join(format!("{artifact_id}.{extension}")))
265    }
266}
267
268pub fn sha256_hex(bytes: &[u8]) -> String {
269    format!("{:x}", Sha256::digest(bytes))
270}
271
272fn validate_artifact_id(artifact_id: &str) -> Result<(), ArtifactStoreError> {
273    if artifact_id.is_empty()
274        || artifact_id.len() > 128
275        || !artifact_id
276            .bytes()
277            .all(|byte| byte.is_ascii_alphanumeric() || byte == b'_' || byte == b'-')
278    {
279        return Err(ArtifactStoreError::InvalidId);
280    }
281    Ok(())
282}