1use 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#[derive(Debug, Clone, PartialEq, Eq)]
22pub struct ArtifactChunk {
23 pub offset: u64,
24 pub bytes: Vec<u8>,
25 pub eof: bool,
26}
27
28#[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#[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}