kcode_pending_object_store/
lib.rs1use std::{
8 collections::HashSet,
9 fs::{File, OpenOptions},
10 io::{Read, Seek, Write},
11 path::{Path, PathBuf},
12};
13
14use anyhow::{Context as _, ensure};
15use sha2::{Digest, Sha256};
16
17const MAGIC: &[u8] = b"KSPENDING01\n";
18const FIXED_HEADER_BYTES: usize = 2 + 4 + 4 + 8 + 32;
19
20#[derive(Clone, Debug, Eq, PartialEq)]
21pub struct StoredPendingObject {
22 pub file_name: String,
23 pub media_type: String,
24 pub bytes: Vec<u8>,
25}
26
27#[derive(Clone, Debug)]
28pub struct PendingObjectStore {
29 directory: PathBuf,
30 namespace: String,
31 format_id: String,
32}
33
34impl PendingObjectStore {
35 pub fn new(
36 directory: impl Into<PathBuf>,
37 namespace: impl Into<String>,
38 format_id: impl Into<String>,
39 ) -> anyhow::Result<Self> {
40 let namespace = namespace.into();
41 let format_id = format_id.into();
42 validate_namespace(&namespace)?;
43 ensure!(!format_id.is_empty(), "format ID cannot be empty");
44 u16::try_from(format_id.len()).context("format ID is too long")?;
45 Ok(Self {
46 directory: directory.into(),
47 namespace,
48 format_id,
49 })
50 }
51
52 pub fn install(
53 &self,
54 position: u64,
55 file_name: &str,
56 media_type: &str,
57 bytes: &[u8],
58 ) -> anyhow::Result<()> {
59 ensure!(
60 !file_name.trim().is_empty(),
61 "object filename cannot be empty"
62 );
63 ensure!(
64 !media_type.trim().is_empty(),
65 "object media type cannot be empty"
66 );
67
68 let format_id = self.format_id.as_bytes();
69 let format_id_len = u16::try_from(format_id.len()).context("format ID is too long")?;
70 let file_name_len =
71 u32::try_from(file_name.len()).context("object filename exceeds 4 GiB")?;
72 let media_type_len =
73 u32::try_from(media_type.len()).context("object media type exceeds 4 GiB")?;
74 let object_len = u64::try_from(bytes.len()).context("object exceeds addressable size")?;
75 let final_path = self.final_path(position);
76 let temp_path = self.temp_path(position);
77
78 ensure!(
79 !final_path.exists(),
80 "pending object file {} already exists",
81 final_path.display()
82 );
83 if temp_path.exists() {
84 std::fs::remove_file(&temp_path)
85 .with_context(|| format!("removing stale temporary {}", temp_path.display()))?;
86 }
87
88 let mut file = OpenOptions::new()
89 .create_new(true)
90 .write(true)
91 .open(&temp_path)
92 .with_context(|| format!("creating temporary {}", temp_path.display()))?;
93 file.write_all(MAGIC)?;
94 file.write_all(&format_id_len.to_le_bytes())?;
95 file.write_all(&file_name_len.to_le_bytes())?;
96 file.write_all(&media_type_len.to_le_bytes())?;
97 file.write_all(&object_len.to_le_bytes())?;
98 file.write_all(&Sha256::digest(bytes))?;
99 file.write_all(format_id)?;
100 file.write_all(file_name.as_bytes())?;
101 file.write_all(media_type.as_bytes())?;
102 file.write_all(bytes)?;
103 file.sync_all()
104 .with_context(|| format!("synchronizing temporary {}", temp_path.display()))?;
105 std::fs::rename(&temp_path, &final_path).with_context(|| {
106 format!(
107 "renaming {} to {}",
108 temp_path.display(),
109 final_path.display()
110 )
111 })?;
112 sync_directory(&self.directory)
113 }
114
115 pub fn read(&self, position: u64) -> anyhow::Result<StoredPendingObject> {
116 let path = self.final_path(position);
117 let mut file = File::open(&path)
118 .with_context(|| format!("opening pending object {}", path.display()))?;
119
120 let mut magic = vec![0_u8; MAGIC.len()];
121 file.read_exact(&mut magic)?;
122 ensure!(
123 magic == MAGIC,
124 "{} is not a pending-object file",
125 path.display()
126 );
127
128 let mut fixed = [0_u8; FIXED_HEADER_BYTES];
129 file.read_exact(&mut fixed)?;
130 let format_id_len = u16::from_le_bytes(fixed[0..2].try_into().unwrap()) as usize;
131 let file_name_len = u32::from_le_bytes(fixed[2..6].try_into().unwrap()) as usize;
132 let media_type_len = u32::from_le_bytes(fixed[6..10].try_into().unwrap()) as usize;
133 let object_len = u64::from_le_bytes(fixed[10..18].try_into().unwrap());
134 let checksum = &fixed[18..50];
135
136 let variable_len = format_id_len
137 .checked_add(file_name_len)
138 .and_then(|value| value.checked_add(media_type_len))
139 .context("pending-object header length overflow")?;
140 let expected_file_len = u64::try_from(MAGIC.len() + FIXED_HEADER_BYTES)
141 .context("pending-object fixed header does not fit u64")?
142 .checked_add(
143 u64::try_from(variable_len)
144 .context("pending-object variable header does not fit u64")?,
145 )
146 .and_then(|value| value.checked_add(object_len))
147 .context("pending-object declared length overflow")?;
148 ensure!(
149 file.metadata()?.len() == expected_file_len,
150 "pending-object declared length differs from file length"
151 );
152
153 let mut variable = vec![0_u8; variable_len];
154 file.read_exact(&mut variable)?;
155 let format_id = std::str::from_utf8(&variable[..format_id_len])?;
156 ensure!(
157 format_id == self.format_id,
158 "unsupported pending-object format {format_id}"
159 );
160
161 let file_name_end = format_id_len + file_name_len;
162 let file_name = std::str::from_utf8(&variable[format_id_len..file_name_end])?.to_owned();
163 let media_type =
164 std::str::from_utf8(&variable[file_name_end..file_name_end + media_type_len])?
165 .to_owned();
166 ensure!(!file_name.trim().is_empty(), "object filename is empty");
167 ensure!(!media_type.trim().is_empty(), "object media type is empty");
168
169 let object_len =
170 usize::try_from(object_len).context("pending object does not fit memory")?;
171 let mut bytes = vec![0_u8; object_len];
172 file.read_exact(&mut bytes)?;
173 ensure!(
174 file.stream_position()? == expected_file_len,
175 "pending-object file has trailing bytes"
176 );
177 let actual_checksum = Sha256::digest(&bytes);
178 ensure!(
179 &actual_checksum[..] == checksum,
180 "pending-object checksum mismatch"
181 );
182
183 Ok(StoredPendingObject {
184 file_name,
185 media_type,
186 bytes,
187 })
188 }
189
190 pub fn verify_all(&self, referenced_positions: &[u64]) -> anyhow::Result<()> {
191 for &position in referenced_positions {
192 self.read(position)?;
193 }
194 Ok(())
195 }
196
197 pub fn reconcile(&self, referenced_positions: &[u64]) -> anyhow::Result<()> {
198 let referenced = referenced_positions.iter().copied().collect::<HashSet<_>>();
199 let mut removed = false;
200
201 if self.directory.exists() {
202 for entry in std::fs::read_dir(&self.directory)
203 .with_context(|| format!("reading directory {}", self.directory.display()))?
204 {
205 let entry = entry?;
206 let Some(name) = entry.file_name().to_str().map(str::to_owned) else {
207 continue;
208 };
209 let Some((position, temporary)) = self.parse_file_name(&name) else {
210 continue;
211 };
212 if temporary || !referenced.contains(&position) {
213 std::fs::remove_file(entry.path())
214 .with_context(|| format!("removing {}", entry.path().display()))?;
215 removed = true;
216 }
217 }
218 }
219
220 if removed {
221 sync_directory(&self.directory)?;
222 }
223 self.verify_all(referenced_positions)
224 }
225
226 pub fn delete_all(&self) -> anyhow::Result<()> {
227 if self.directory.exists() {
228 for entry in std::fs::read_dir(&self.directory)
229 .with_context(|| format!("reading directory {}", self.directory.display()))?
230 {
231 let entry = entry?;
232 let Some(name) = entry.file_name().to_str().map(str::to_owned) else {
233 continue;
234 };
235 if self.parse_file_name(&name).is_some() {
236 std::fs::remove_file(entry.path())
237 .with_context(|| format!("removing {}", entry.path().display()))?;
238 }
239 }
240 }
241 sync_directory(&self.directory)
242 }
243
244 fn final_path(&self, position: u64) -> PathBuf {
245 self.directory
246 .join(format!("{}-{position}.pending-object", self.namespace))
247 }
248
249 fn temp_path(&self, position: u64) -> PathBuf {
250 self.directory
251 .join(format!("{}-{position}.pending-object.tmp", self.namespace))
252 }
253
254 fn parse_file_name(&self, name: &str) -> Option<(u64, bool)> {
255 let tail = name.strip_prefix(&format!("{}-", self.namespace))?;
256 let (number, temporary) = if let Some(number) = tail.strip_suffix(".pending-object.tmp") {
257 (number, true)
258 } else {
259 (tail.strip_suffix(".pending-object")?, false)
260 };
261 if number.is_empty() || number.starts_with('0') && number != "0" {
262 return None;
263 }
264 let position = number.parse::<u64>().ok()?;
265 (position.to_string() == number).then_some((position, temporary))
266 }
267}
268
269fn validate_namespace(namespace: &str) -> anyhow::Result<()> {
270 ensure!(!namespace.is_empty(), "namespace cannot be empty");
271 ensure!(namespace.len() <= 255, "namespace exceeds 255 characters");
272 ensure!(
273 namespace
274 .bytes()
275 .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_')),
276 "namespace contains characters that are unsafe in filenames"
277 );
278 Ok(())
279}
280
281fn sync_directory(directory: &Path) -> anyhow::Result<()> {
282 File::open(directory)
283 .with_context(|| format!("opening directory {}", directory.display()))?
284 .sync_all()
285 .with_context(|| format!("synchronizing directory {}", directory.display()))
286}
287
288#[cfg(test)]
289mod tests {
290 use std::time::{SystemTime, UNIX_EPOCH};
291
292 use super::*;
293
294 struct TestDirectory(PathBuf);
295
296 impl TestDirectory {
297 fn new(label: &str) -> Self {
298 Self(std::env::temp_dir().join(format!(
299 "pending-object-store-{label}-{}-{}",
300 std::process::id(),
301 SystemTime::now()
302 .duration_since(UNIX_EPOCH)
303 .unwrap()
304 .as_nanos()
305 )))
306 }
307 }
308
309 impl Drop for TestDirectory {
310 fn drop(&mut self) {
311 if self.0.exists() {
312 std::fs::remove_dir_all(&self.0).unwrap();
313 }
314 }
315 }
316
317 #[test]
318 fn construction_validates_configuration_without_touching_filesystem() {
319 let directory = TestDirectory::new("new");
320 let store = PendingObjectStore::new(&directory.0, "session_1", "0.2.1").unwrap();
321 assert!(!store.directory.exists());
322 assert!(PendingObjectStore::new(&directory.0, "", "0.2.1").is_err());
323 assert!(PendingObjectStore::new(&directory.0, "bad/name", "0.2.1").is_err());
324 assert!(PendingObjectStore::new(&directory.0, "a".repeat(256), "0.2.1").is_err());
325 assert!(PendingObjectStore::new(&directory.0, "valid", "").is_err());
326 assert!(PendingObjectStore::new(&directory.0, "valid", "x".repeat(65_536)).is_err());
327 assert!(!directory.0.exists());
328 }
329
330 #[test]
331 fn install_matches_the_session_log_format_and_read_verifies_it() {
332 let directory = TestDirectory::new("format");
333 std::fs::create_dir(&directory.0).unwrap();
334 let store = PendingObjectStore::new(&directory.0, "session-7", "0.2.1").unwrap();
335 store
336 .install(3, "notes.txt", "text/plain", b"durable bytes")
337 .unwrap();
338
339 let mut expected = Vec::new();
340 expected.extend_from_slice(MAGIC);
341 expected.extend_from_slice(&5_u16.to_le_bytes());
342 expected.extend_from_slice(&9_u32.to_le_bytes());
343 expected.extend_from_slice(&10_u32.to_le_bytes());
344 expected.extend_from_slice(&13_u64.to_le_bytes());
345 expected.extend_from_slice(&Sha256::digest(b"durable bytes"));
346 expected.extend_from_slice(b"0.2.1notes.txttext/plaindurable bytes");
347 assert_eq!(
348 std::fs::read(directory.0.join("session-7-3.pending-object")).unwrap(),
349 expected
350 );
351 assert!(!directory.0.join("session-7-3.pending-object.tmp").exists());
352 assert_eq!(
353 store.read(3).unwrap(),
354 StoredPendingObject {
355 file_name: "notes.txt".into(),
356 media_type: "text/plain".into(),
357 bytes: b"durable bytes".to_vec(),
358 }
359 );
360
361 assert!(
362 store
363 .install(3, "replacement", "text/plain", b"new")
364 .is_err()
365 );
366 assert_eq!(
367 std::fs::read(directory.0.join("session-7-3.pending-object")).unwrap(),
368 expected
369 );
370 }
371
372 #[test]
373 fn install_removes_only_its_exact_stale_temporary_file() {
374 let directory = TestDirectory::new("temporary");
375 std::fs::create_dir(&directory.0).unwrap();
376 let store = PendingObjectStore::new(&directory.0, "space", "format").unwrap();
377 let stale = directory.0.join("space-4.pending-object.tmp");
378 let other = directory.0.join("space-04.pending-object.tmp");
379 std::fs::write(&stale, b"stale").unwrap();
380 std::fs::write(&other, b"keep").unwrap();
381
382 store.install(4, "a", "b", b"c").unwrap();
383
384 assert!(!stale.exists());
385 assert!(other.exists());
386 assert!(
387 PendingObjectStore::new(directory.0.join("missing"), "space", "format")
388 .unwrap()
389 .install(0, "a", "b", b"c")
390 .is_err()
391 );
392 }
393
394 #[test]
395 fn reads_reject_wrong_format_length_checksum_and_empty_metadata() {
396 let directory = TestDirectory::new("invalid");
397 std::fs::create_dir(&directory.0).unwrap();
398 let store = PendingObjectStore::new(&directory.0, "object", "format-a").unwrap();
399
400 store.install(0, "file", "type", b"bytes").unwrap();
401 assert!(
402 PendingObjectStore::new(&directory.0, "object", "format-b")
403 .unwrap()
404 .read(0)
405 .is_err()
406 );
407
408 let path = directory.0.join("object-0.pending-object");
409 let valid = std::fs::read(&path).unwrap();
410 let mut corrupted = valid.clone();
411 *corrupted.last_mut().unwrap() ^= 0xff;
412 std::fs::write(&path, &corrupted).unwrap();
413 assert!(store.read(0).is_err());
414
415 let mut trailing = valid.clone();
416 trailing.push(0);
417 std::fs::write(&path, trailing).unwrap();
418 assert!(store.read(0).is_err());
419
420 let mut empty_name = valid;
421 let fixed_start = MAGIC.len();
422 empty_name[fixed_start + 2..fixed_start + 6].copy_from_slice(&0_u32.to_le_bytes());
423 std::fs::write(&path, empty_name).unwrap();
424 assert!(store.read(0).is_err());
425 }
426
427 #[test]
428 fn verification_does_not_clean_and_reconciliation_is_exact() {
429 let directory = TestDirectory::new("reconcile");
430 std::fs::create_dir(&directory.0).unwrap();
431 let store = PendingObjectStore::new(&directory.0, "abc", "format").unwrap();
432 store.install(1, "one", "type", b"one").unwrap();
433 store.install(2, "two", "type", b"two").unwrap();
434
435 let temporary = directory.0.join("abc-8.pending-object.tmp");
436 let leading_zero = directory.0.join("abc-08.pending-object.tmp");
437 let similar_prefix = directory.0.join("abc-other-9.pending-object");
438 let malformed = directory.0.join("abc-x.pending-object");
439 std::fs::write(&temporary, b"remove").unwrap();
440 std::fs::write(&leading_zero, b"keep").unwrap();
441 std::fs::write(&similar_prefix, b"keep").unwrap();
442 std::fs::write(&malformed, b"keep").unwrap();
443
444 store.verify_all(&[1]).unwrap();
445 assert!(temporary.exists());
446
447 store.reconcile(&[1]).unwrap();
448 assert!(directory.0.join("abc-1.pending-object").exists());
449 assert!(!directory.0.join("abc-2.pending-object").exists());
450 assert!(!temporary.exists());
451 assert!(leading_zero.exists());
452 assert!(similar_prefix.exists());
453 assert!(malformed.exists());
454 assert!(store.reconcile(&[99]).is_err());
455 }
456
457 #[test]
458 fn delete_all_removes_only_recognized_names() {
459 let directory = TestDirectory::new("delete");
460 std::fs::create_dir(&directory.0).unwrap();
461 let store = PendingObjectStore::new(&directory.0, "abc", "format").unwrap();
462 store.install(0, "zero", "type", b"zero").unwrap();
463
464 let temporary = directory.0.join("abc-1.pending-object.tmp");
465 let leading_zero = directory.0.join("abc-01.pending-object");
466 let similar_prefix = directory.0.join("abcd-2.pending-object");
467 std::fs::write(&temporary, b"remove").unwrap();
468 std::fs::write(&leading_zero, b"keep").unwrap();
469 std::fs::write(&similar_prefix, b"keep").unwrap();
470
471 store.delete_all().unwrap();
472
473 assert!(!directory.0.join("abc-0.pending-object").exists());
474 assert!(!temporary.exists());
475 assert!(leading_zero.exists());
476 assert!(similar_prefix.exists());
477 }
478}