Skip to main content

objectstore_service/backend/
local_fs.rs

1//! Local filesystem backend for development and testing.
2//!
3//! Complete object files are published by atomically renaming same-directory drafts, so readers
4//! observe either the previous or next complete file. Unpublished drafts are removed automatically
5//! when dropped.
6//!
7//! To avoid races on metadata, expiry, and upload updates, this backend uses locks placed under
8//! `.locks/` to synchronize mutations across backend instances and cooperating processes. The first
9//! two bytes of the BLAKE3 hash of an object's storage path select one of 65,536 permanent lock
10//! slots under `.locks/<first byte>/<second byte>` (lowercase hexadecimal). An object and all of its
11//! resumable uploads share the same lock.
12//!
13//! Shared filesystems are supported only when locks propagate across the cluster, pathname
14//! visibility is coherent, and rename is atomic.
15//!
16//! Newly created files and directories are owner-only on Unix (0600 and 0700 respectively,
17//! further restricted by umask). Existing paths retain their permissions.
18
19use std::fs::File;
20use std::io;
21use std::num::NonZeroU64;
22use std::path::{Path, PathBuf};
23use std::pin::pin;
24use std::sync::Arc;
25use std::time::SystemTime;
26
27use futures_util::StreamExt;
28use objectstore_types::metadata::Metadata;
29use objectstore_types::range::ByteRange;
30use objectstore_types::resumable::UploadProgress;
31use objectstore_types::time::Timestamp;
32use tokio::fs::OpenOptions;
33use tokio::io::{
34    AsyncBufReadExt, AsyncRead, AsyncReadExt, AsyncSeekExt, AsyncWriteExt, BufReader, BufWriter,
35};
36use tokio::sync::Semaphore;
37use tokio_util::io::{ReaderStream, StreamReader};
38use uuid::Uuid;
39
40use crate::backend::common::{
41    Backend, DeleteResponse, GetResponse, MultipartUploadBackend, PutResponse,
42};
43use crate::change_stream::{
44    ChangeStream, ChangeStreamFactory, CostTrackerStreamConfig, flush_change_stream,
45};
46use crate::error::{Error, ErrorKind, Result, ResultExt as _};
47use crate::id::ObjectId;
48use crate::multipart::{
49    AbortMultipartResponse, CompleteMultipartResponse, CompletedPart, InitiateMultipartResponse,
50    ListPartsResponse, Part, PartNumber, UploadId, UploadPartResponse,
51};
52use crate::resumable::BackendToken;
53use crate::stream::{self, ClientStream};
54
55/// Options for owner-only files on Unix, further restricted by umask.
56fn file_options() -> std::fs::OpenOptions {
57    let mut options = std::fs::OpenOptions::new();
58    #[cfg(unix)]
59    {
60        use std::os::unix::fs::OpenOptionsExt as _;
61        options.mode(0o600);
62    }
63    options
64}
65
66/// Creates owner-only directories on Unix without changing existing permissions.
67async fn create_directories(path: &Path) -> io::Result<()> {
68    let mut builder = tokio::fs::DirBuilder::new();
69    builder.recursive(true);
70    #[cfg(unix)]
71    builder.mode(0o700);
72    builder.create(path).await
73}
74
75/// Configuration for [`LocalFsBackend`].
76///
77/// Stores objects as files on the local filesystem. Suitable for development, testing,
78/// and single-server deployments.
79///
80/// # Example
81///
82/// ```yaml
83/// storage:
84///   type: filesystem
85///   path: /data
86/// ```
87#[derive(Debug, Clone, serde::Deserialize, serde::Serialize)]
88pub struct FileSystemConfig {
89    /// Directory path for storing objects.
90    ///
91    /// The directory will be created if it doesn't exist. Relative paths are resolved from
92    /// the server's working directory.
93    ///
94    /// # Default
95    ///
96    /// `"data"` (relative to the server's working directory)
97    ///
98    /// # Environment Variables
99    ///
100    /// - `OS__STORAGE__TYPE=filesystem`
101    /// - `OS__STORAGE__PATH=/path/to/storage`
102    pub path: PathBuf,
103
104    /// Reports what this backend stores, for per-usecase cost attribution.
105    ///
106    /// # Default
107    ///
108    /// `None`, which disables reporting for this backend.
109    ///
110    /// # Environment Variables
111    ///
112    /// - `OS__STORAGE__COGS__SHARED_RESOURCE_ID=filesystem_objectstore`
113    /// - `OS__STORAGE__COGS__SAMPLE_RATE=1.0` (optional)
114    #[serde(default, skip_serializing_if = "Option::is_none")]
115    pub cogs: Option<CostTrackerStreamConfig>,
116}
117
118/// Local filesystem backend for development and testing.
119#[derive(Debug)]
120pub struct LocalFsBackend {
121    path: PathBuf,
122    locks: ObjectLocks,
123
124    change_stream: Arc<dyn ChangeStream>,
125}
126
127impl LocalFsBackend {
128    /// Creates a new [`LocalFsBackend`] rooted at the directory in `config`.
129    pub fn new(config: FileSystemConfig, streams: &ChangeStreamFactory) -> Self {
130        let FileSystemConfig { path, cogs } = config;
131        let locks = ObjectLocks::new(&path);
132        Self {
133            path,
134            locks,
135            change_stream: streams.build(cogs.as_ref()),
136        }
137    }
138
139    /// Returns the filesystem path for the given object ID.
140    fn path(&self, id: &ObjectId) -> PathBuf {
141        self.path.join(id.as_storage_path().to_string())
142    }
143
144    fn upload_path(&self, upload_id: Uuid) -> PathBuf {
145        self.path.join("uploads").join(upload_id.to_string())
146    }
147}
148
149#[async_trait::async_trait]
150impl Backend for LocalFsBackend {
151    fn name(&self) -> &'static str {
152        "local-fs"
153    }
154
155    fn as_multipart_upload_backend(&self) -> Result<&dyn MultipartUploadBackend> {
156        Ok(self)
157    }
158
159    #[tracing::instrument(level = "debug", fields(?id), skip_all)]
160    async fn put_object(
161        &self,
162        id: &ObjectId,
163        metadata: &Metadata,
164        stream: ClientStream,
165        _access_time: Timestamp,
166    ) -> Result<PutResponse> {
167        let path = self.path(id);
168        objectstore_log::debug!(path=%path.display(), "Writing to local_fs backend");
169        create_directories(path.parent().unwrap()).await.context(
170            ErrorKind::BackendFailure,
171            "creating local-fs object directory",
172        )?;
173
174        let mut draft = Draft::create(&path, metadata).await?;
175        let mut reader = pin!(StreamReader::new(stream));
176        let payload_size = tokio::io::copy(&mut reader, draft.writer())
177            .await
178            .map_err(|e| match stream::unpack_client_error(&e) {
179                Some(ce) => Error::from(ce),
180                None => Error::with_context(
181                    ErrorKind::BackendFailure,
182                    "writing local-fs object payload",
183                    e,
184                ),
185            })?;
186
187        let stored_size = draft.preamble_len() + payload_size;
188
189        draft.prepare().await?;
190        let _guard = self.locks.acquire(id).await?;
191        draft.publish().await?;
192
193        self.change_stream
194            .write(id, stored_size, metadata.time_expires);
195
196        Ok(())
197    }
198
199    #[tracing::instrument(level = "debug", skip(self))]
200    async fn get_object(
201        &self,
202        id: &ObjectId,
203        access_time: Timestamp,
204        range: Option<ByteRange>,
205    ) -> Result<GetResponse> {
206        objectstore_log::debug!("Reading from local_fs backend");
207        let path = self.path(id);
208        let Some(object) = ObjectFile::try_open(&path, access_time).await? else {
209            objectstore_log::debug!("Object not found");
210            return Ok(None);
211        };
212        let ObjectFile {
213            metadata,
214            preamble_len,
215            payload_size,
216            mut reader,
217        } = object;
218
219        let (content_range, stream) = match range {
220            Some(byte_range) => {
221                let content_range =
222                    byte_range
223                        .resolve(payload_size)
224                        .ok_or(ErrorKind::RangeNotSatisfiable {
225                            total: payload_size,
226                        })?;
227                let payload_start = preamble_len + content_range.start;
228                reader
229                    .seek(std::io::SeekFrom::Start(payload_start))
230                    .await
231                    .context(ErrorKind::BackendFailure, "seeking local-fs object payload")?;
232                let limited = reader.take(content_range.len());
233                (Some(content_range), ReaderStream::new(limited).boxed())
234            }
235            None => (None, ReaderStream::new(reader).boxed()),
236        };
237        Ok(Some((metadata, content_range, stream)))
238    }
239
240    #[tracing::instrument(level = "debug", skip(self))]
241    async fn set_expiry(
242        &self,
243        id: &ObjectId,
244        expire_at: Timestamp,
245        access_time: Timestamp,
246    ) -> Result<bool> {
247        let _guard = self.locks.acquire(id).await?;
248
249        let path = self.path(id);
250        let Some(object) = ObjectFile::try_open(&path, access_time).await? else {
251            return Ok(false);
252        };
253        let ObjectFile {
254            mut metadata,
255            mut reader,
256            ..
257        } = object;
258
259        let Some(current_expiry) = metadata.time_expires else {
260            return Ok(false);
261        };
262        if current_expiry >= expire_at {
263            return Ok(true); // already satisfied
264        }
265        metadata.time_expires = Some(expire_at);
266
267        let mut draft = Draft::create(&path, &metadata).await?;
268        tokio::io::copy(&mut reader, draft.writer()).await.context(
269            ErrorKind::BackendFailure,
270            "copying local-fs object payload for expiry extension",
271        )?;
272
273        draft.prepare().await?;
274        draft.publish().await?;
275
276        self.change_stream.update(id, Some(expire_at));
277
278        Ok(true)
279    }
280
281    #[tracing::instrument(level = "debug", skip(self))]
282    async fn delete_object(
283        &self,
284        id: &ObjectId,
285        _access_time: Timestamp,
286    ) -> Result<DeleteResponse> {
287        let _guard = self.locks.acquire(id).await?;
288
289        objectstore_log::debug!("Deleting from local_fs backend");
290        let path = self.path(id);
291        match tokio::fs::remove_file(path).await {
292            Ok(()) => self.change_stream.delete(id),
293            Err(error) if error.kind() == io::ErrorKind::NotFound => {
294                objectstore_log::debug!("Object not found");
295            }
296            result => {
297                result.context(ErrorKind::BackendFailure, "deleting local-fs object")?;
298            }
299        }
300
301        Ok(())
302    }
303
304    #[tracing::instrument(level = "debug", fields(?id, total_length), skip_all)]
305    async fn create_upload_session(
306        &self,
307        id: &ObjectId,
308        metadata: &Metadata,
309        total_length: NonZeroU64,
310    ) -> Result<Option<BackendToken>> {
311        let upload_id = uuid::Uuid::now_v7();
312        let path = self.upload_path(upload_id);
313        create_directories(path.parent().unwrap()).await.context(
314            ErrorKind::BackendFailure,
315            "creating local-fs object directory",
316        )?;
317        UploadFile::create(&path, metadata).await?;
318        Ok(Some(format!("{total_length}.{upload_id}")))
319    }
320
321    // In this backend, if all the bytes of a resumable upload have been written but publication failed,
322    // an empty chunk or offset query will retry publication.
323    #[tracing::instrument(level = "debug", fields(?id, offset, content_length), skip_all)]
324    async fn put_chunk(
325        &self,
326        id: &ObjectId,
327        token: &BackendToken,
328        offset: u64,
329        content_length: u64,
330        stream: ClientStream,
331    ) -> Result<UploadProgress> {
332        let session = UploadSession::from_token(token)?;
333        offset
334            .checked_add(content_length)
335            .filter(|end| *end <= session.total_length.get())
336            .ok_or(ErrorKind::ChunkExceedsUploadLength {
337                offset,
338                content_length,
339                upload_length: session.total_length.get(),
340            })?;
341        let _guard = self.locks.acquire(id).await?;
342
343        let upload_path = self.upload_path(session.upload_id);
344        let mut upload = UploadFile::open(&upload_path).await?;
345        if content_length != 0 && offset != upload.offset() {
346            return Err(ErrorKind::UploadOffsetMismatch {
347                offset: upload.offset(),
348            }
349            .into());
350        }
351
352        let reader = StreamReader::new(stream).take(content_length);
353        let persisted_offset = upload.append(reader).await?;
354
355        if persisted_offset != session.total_length.get() {
356            return Ok(UploadProgress::Incomplete {
357                offset: persisted_offset,
358            });
359        }
360
361        let object_path = self.path(id);
362        create_directories(object_path.parent().unwrap())
363            .await
364            .context(
365                ErrorKind::BackendFailure,
366                "creating local-fs object directory",
367            )?;
368        let (stored_size, expires_at) = upload.publish(object_path).await?;
369        self.change_stream.write(id, stored_size, expires_at);
370        Ok(UploadProgress::Complete)
371    }
372
373    #[tracing::instrument(level = "debug", fields(?id), skip_all)]
374    async fn upload_offset(&self, id: &ObjectId, token: &BackendToken) -> Result<UploadProgress> {
375        self.put_chunk(id, token, 0, 0, futures_util::stream::empty().boxed())
376            .await
377    }
378
379    #[tracing::instrument(level = "debug", fields(?id), skip_all)]
380    async fn cancel_upload(&self, id: &ObjectId, token: &BackendToken) -> Result<()> {
381        let session = UploadSession::from_token(token)?;
382        let _guard = self.locks.acquire(id).await?;
383        let path = self.upload_path(session.upload_id);
384        match tokio::fs::remove_file(&path).await {
385            Ok(()) => Ok(()),
386            Err(error) if error.kind() == io::ErrorKind::NotFound => {
387                Err(ErrorKind::UnknownUploadSession.into())
388            }
389            result => result.context(
390                ErrorKind::BackendFailure,
391                "canceling local-fs resumable upload",
392            ),
393        }
394    }
395
396    async fn join(&self) {
397        flush_change_stream(&self.change_stream).await;
398    }
399}
400
401impl LocalFsBackend {
402    fn multipart_dir(&self, id: &ObjectId, upload_id: &UploadId) -> PathBuf {
403        self.path
404            .join("__multipart__")
405            .join(id.as_storage_path().to_string())
406            .join(upload_id.as_str())
407    }
408}
409
410#[async_trait::async_trait]
411impl MultipartUploadBackend for LocalFsBackend {
412    async fn initiate_multipart(
413        &self,
414        id: &ObjectId,
415        metadata: &Metadata,
416    ) -> Result<InitiateMultipartResponse> {
417        let upload_id = UploadId::new(Uuid::now_v7().to_string())?;
418        let dir = self.multipart_dir(id, &upload_id);
419        create_directories(&dir).await.context(
420            ErrorKind::BackendFailure,
421            "creating local-fs multipart upload",
422        )?;
423
424        let meta_path = dir.join("metadata.json");
425        let metadata_json = serde_json::to_string(metadata)
426            .context(ErrorKind::Internal, "encoding local-fs multipart metadata")?;
427        let mut file = OpenOptions::from(file_options())
428            .create(true)
429            .truncate(true)
430            .write(true)
431            .open(meta_path)
432            .await
433            .context(
434                ErrorKind::BackendFailure,
435                "creating local-fs multipart metadata",
436            )?;
437        file.write_all(metadata_json.as_bytes()).await.context(
438            ErrorKind::BackendFailure,
439            "writing local-fs multipart metadata",
440        )?;
441        file.flush().await.context(
442            ErrorKind::BackendFailure,
443            "flushing local-fs multipart metadata",
444        )?;
445
446        Ok(upload_id)
447    }
448
449    async fn upload_part(
450        &self,
451        id: &ObjectId,
452        upload_id: &UploadId,
453        part_number: PartNumber,
454        content_length: u64,
455        _content_md5: Option<&str>,
456        body: ClientStream,
457    ) -> Result<UploadPartResponse> {
458        let dir = self.multipart_dir(id, upload_id);
459        if !tokio::fs::try_exists(&dir).await.context(
460            ErrorKind::BackendFailure,
461            "checking local-fs multipart upload",
462        )? {
463            return Err(Error::new(
464                ErrorKind::BackendFailure,
465                "local-fs multipart upload not found",
466            ));
467        }
468
469        let etag = format!("\"etag-{part_number}-{content_length}\"");
470
471        let header = serde_json::json!({
472            "etag": etag,
473            "uploaded_at": SystemTime::now(),
474            "size": content_length,
475        });
476        let header_line = serde_json::to_string(&header)
477            .context(ErrorKind::Internal, "encoding local-fs part header")?;
478
479        let part_path = dir.join(format!("{part_number}.part"));
480        let file = OpenOptions::from(file_options())
481            .create(true)
482            .write(true)
483            .truncate(true)
484            .open(part_path)
485            .await
486            .context(
487                ErrorKind::BackendFailure,
488                "opening local-fs multipart part for writing",
489            )?;
490
491        let mut reader = pin!(StreamReader::new(body));
492        let mut writer = BufWriter::new(file);
493        writer
494            .write_all(header_line.as_bytes())
495            .await
496            .context(ErrorKind::BackendFailure, "writing local-fs part header")?;
497        writer
498            .write_all(b"\n")
499            .await
500            .context(ErrorKind::BackendFailure, "writing local-fs part header")?;
501
502        let _bytes_copied = tokio::io::copy(&mut reader, &mut writer)
503            .await
504            .map_err(|e| match stream::unpack_client_error(&e) {
505                Some(ce) => Error::from(ce),
506                None => Error::with_context(
507                    ErrorKind::BackendFailure,
508                    "writing local-fs multipart part payload",
509                    e,
510                ),
511            })?;
512
513        // TODO: validate bytes_copied against content_length and return a BadRequest-style
514        // error. Needs a service-layer error variant that maps to HTTP 400 without abusing
515        // ClientError (which is meant for stream errors).
516
517        writer.flush().await.context(
518            ErrorKind::BackendFailure,
519            "flushing local-fs multipart part",
520        )?;
521        let file = writer.into_inner();
522        file.sync_data()
523            .await
524            .context(ErrorKind::BackendFailure, "syncing local-fs multipart part")?;
525        drop(file);
526
527        Ok(etag)
528    }
529
530    async fn list_parts(
531        &self,
532        id: &ObjectId,
533        upload_id: &UploadId,
534        max_parts: Option<u32>,
535        part_number_marker: Option<PartNumber>,
536    ) -> Result<ListPartsResponse> {
537        let dir = self.multipart_dir(id, upload_id);
538        if !tokio::fs::try_exists(&dir).await.context(
539            ErrorKind::BackendFailure,
540            "checking local-fs multipart upload",
541        )? {
542            return Err(Error::new(
543                ErrorKind::BackendFailure,
544                "local-fs multipart upload not found",
545            ));
546        }
547
548        let mut entries = tokio::fs::read_dir(&dir).await.context(
549            ErrorKind::BackendFailure,
550            "listing local-fs multipart parts",
551        )?;
552        let mut parts = Vec::new();
553
554        while let Some(entry) = entries.next_entry().await.context(
555            ErrorKind::BackendFailure,
556            "listing local-fs multipart parts",
557        )? {
558            let name = entry.file_name();
559            let name_str = name.to_string_lossy();
560            let Some(pn_str) = name_str.strip_suffix(".part") else {
561                continue;
562            };
563            let Ok(pn) = pn_str.parse::<PartNumber>() else {
564                continue;
565            };
566
567            if part_number_marker.is_some_and(|marker| pn <= marker) {
568                continue;
569            }
570
571            let file = tokio::fs::File::open(entry.path())
572                .await
573                .context(ErrorKind::BackendFailure, "opening local-fs multipart part")?;
574            let mut reader = BufReader::new(file);
575            let mut header_line = String::new();
576            reader
577                .read_line(&mut header_line)
578                .await
579                .context(ErrorKind::BackendFailure, "reading local-fs part header")?;
580            let header: serde_json::Value = serde_json::from_str(header_line.trim_end())
581                .context(ErrorKind::CorruptData, "decoding local-fs part header")?;
582
583            parts.push(Part {
584                part_number: pn,
585                etag: header["etag"].as_str().unwrap_or("").to_string(),
586                last_modified: serde_json::from_value(header["uploaded_at"].clone())
587                    .unwrap_or(SystemTime::UNIX_EPOCH),
588                size: header["size"].as_u64().unwrap_or(0),
589            });
590        }
591
592        parts.sort_by_key(|p| p.part_number);
593
594        let max = max_parts.unwrap_or(u32::MAX) as usize;
595        let is_truncated = parts.len() > max;
596        parts.truncate(max);
597
598        let next_part_number_marker = if is_truncated {
599            parts.last().map(|p| p.part_number)
600        } else {
601            None
602        };
603
604        Ok(ListPartsResponse {
605            parts,
606            is_truncated,
607            next_part_number_marker,
608        })
609    }
610
611    async fn abort_multipart(
612        &self,
613        id: &ObjectId,
614        upload_id: &UploadId,
615    ) -> Result<AbortMultipartResponse> {
616        let dir = self.multipart_dir(id, upload_id);
617        if tokio::fs::try_exists(&dir).await.context(
618            ErrorKind::BackendFailure,
619            "checking local-fs multipart upload",
620        )? {
621            tokio::fs::remove_dir_all(dir).await.context(
622                ErrorKind::BackendFailure,
623                "removing local-fs multipart upload",
624            )?;
625        }
626        Ok(())
627    }
628
629    async fn complete_multipart(
630        &self,
631        id: &ObjectId,
632        upload_id: &UploadId,
633        parts: Vec<CompletedPart>,
634        _access_time: Timestamp,
635    ) -> Result<CompleteMultipartResponse> {
636        let dir = self.multipart_dir(id, upload_id);
637        if !tokio::fs::try_exists(&dir).await.context(
638            ErrorKind::BackendFailure,
639            "checking local-fs multipart upload",
640        )? {
641            return Err(Error::new(
642                ErrorKind::BackendFailure,
643                "local-fs multipart upload not found",
644            ));
645        }
646
647        // Read metadata
648        let meta_path = dir.join("metadata.json");
649        let meta_bytes = tokio::fs::read(&meta_path).await.context(
650            ErrorKind::BackendFailure,
651            "reading local-fs multipart metadata",
652        )?;
653        let metadata: Metadata = serde_json::from_slice(&meta_bytes).context(
654            ErrorKind::CorruptData,
655            "decoding local-fs multipart metadata",
656        )?;
657
658        // TODO: validate that parts are in ascending part_number order and reject with
659        // InvalidPartOrder if not (matches S3/GCS behavior). Needs a proper client error variant.
660
661        // Validate all parts (headers only) before writing anything
662        for completed in &parts {
663            let part_path = dir.join(format!("{}.part", completed.part_number));
664            if !tokio::fs::try_exists(&part_path).await.context(
665                ErrorKind::BackendFailure,
666                "checking local-fs multipart part",
667            )? {
668                return Ok(Some(crate::multipart::CompleteMultipartError {
669                    code: "InvalidPart".into(),
670                    message: format!("part number {} was not uploaded", completed.part_number),
671                }));
672            }
673
674            let file = tokio::fs::File::open(&part_path)
675                .await
676                .context(ErrorKind::BackendFailure, "opening local-fs multipart part")?;
677            let mut reader = BufReader::new(file);
678            let mut header_line = String::new();
679            reader
680                .read_line(&mut header_line)
681                .await
682                .context(ErrorKind::BackendFailure, "reading local-fs part header")?;
683            let header: serde_json::Value = serde_json::from_str(header_line.trim_end())
684                .context(ErrorKind::CorruptData, "decoding local-fs part header")?;
685
686            let stored_etag = header["etag"].as_str().unwrap_or("");
687            if stored_etag != completed.etag {
688                return Ok(Some(crate::multipart::CompleteMultipartError {
689                    code: "InvalidPart".into(),
690                    message: format!(
691                        "etag mismatch for part {}: expected {}, got {}",
692                        completed.part_number, stored_etag, completed.etag
693                    ),
694                }));
695            }
696        }
697
698        // Assemble the parts into a draft before publishing the object.
699        let path = self.path(id);
700        create_directories(path.parent().unwrap()).await.context(
701            ErrorKind::BackendFailure,
702            "creating local-fs object directory",
703        )?;
704        let mut draft = Draft::create(&path, &metadata).await?;
705
706        let mut payload_size = 0;
707        for completed in &parts {
708            let part_path = dir.join(format!("{}.part", completed.part_number));
709            let file = tokio::fs::File::open(&part_path)
710                .await
711                .context(ErrorKind::BackendFailure, "opening local-fs multipart part")?;
712            let mut reader = BufReader::new(file);
713            let mut header_line = String::new();
714            reader
715                .read_line(&mut header_line)
716                .await
717                .context(ErrorKind::BackendFailure, "reading local-fs part header")?;
718            payload_size += tokio::io::copy(&mut reader, draft.writer()).await.context(
719                ErrorKind::BackendFailure,
720                "assembling local-fs object payload",
721            )?;
722        }
723
724        let stored_size = draft.preamble_len() + payload_size;
725
726        draft.prepare().await?;
727        let guard = self.locks.acquire(id).await?;
728        draft.publish().await?;
729        drop(guard);
730
731        self.change_stream
732            .write(id, stored_size, metadata.time_expires);
733
734        // Clean up multipart state
735        tokio::fs::remove_dir_all(dir).await.context(
736            ErrorKind::BackendFailure,
737            "removing local-fs multipart upload",
738        )?;
739
740        Ok(None)
741    }
742}
743
744// Must be lower than the tokio runtime `max_blocking_threads` setting.
745const MAX_BLOCKING_LOCK_WAITERS: usize = 256;
746
747#[derive(Debug)]
748struct ObjectLocks {
749    root: PathBuf,
750    blocking_waiters: Arc<Semaphore>,
751}
752
753impl ObjectLocks {
754    pub fn new(storage_root: &Path) -> Self {
755        Self {
756            root: storage_root.join(".locks"),
757            blocking_waiters: Arc::new(Semaphore::new(MAX_BLOCKING_LOCK_WAITERS)),
758        }
759    }
760
761    fn path(&self, id: &ObjectId) -> PathBuf {
762        let hash = blake3::hash(id.as_storage_path().to_string().as_bytes());
763        let bytes = hash.as_bytes();
764        self.root
765            .join(format!("{:02x}", bytes[0]))
766            .join(format!("{:02x}", bytes[1]))
767    }
768
769    /// Acquires a slot lock, released when the returned file is dropped.
770    pub async fn acquire(&self, id: &ObjectId) -> Result<File> {
771        let path = self.path(id);
772        create_directories(path.parent().unwrap()).await.context(
773            ErrorKind::BackendFailure,
774            "creating local-fs object lock directory",
775        )?;
776
777        // Leave blocking-pool capacity available for the current lock holder's filesystem work.
778        let permit = Arc::clone(&self.blocking_waiters)
779            .acquire_owned()
780            .await
781            .expect("local-fs lock semaphore is never closed");
782
783        tokio::task::spawn_blocking(move || -> io::Result<File> {
784            let _permit = permit;
785            let file = file_options()
786                .create(true)
787                .truncate(false)
788                .read(true)
789                .write(true)
790                .open(&path)?;
791            file.lock()?;
792            Ok(file)
793        })
794        .await
795        .context(ErrorKind::Internal, "waiting for local-fs object lock")?
796        .context(ErrorKind::BackendFailure, "acquiring local-fs object lock")
797    }
798}
799
800#[derive(Debug)]
801struct UploadSession {
802    total_length: NonZeroU64,
803    upload_id: Uuid,
804}
805
806impl UploadSession {
807    fn from_token(token: &BackendToken) -> Result<Self> {
808        let (length, upload_id) = token
809            .split_once('.')
810            .ok_or(ErrorKind::UnknownUploadSession)?;
811        let total_length = length
812            .parse::<NonZeroU64>()
813            .map_err(|_| ErrorKind::UnknownUploadSession)?;
814        let upload_id_str = upload_id;
815        let upload_id =
816            Uuid::parse_str(upload_id_str).map_err(|_| ErrorKind::UnknownUploadSession)?;
817        Ok(Self {
818            total_length,
819            upload_id,
820        })
821    }
822}
823
824/// An open resumable upload containing a metadata preamble followed by payload bytes.
825struct UploadFile {
826    file: tokio::fs::File,
827    path: PathBuf,
828    /// Number of payload bytes stored after the metadata preamble.
829    payload_size: u64,
830    /// Total size and expiration, for reporting to the `ChangeStream`.
831    stored_size: u64,
832    expires_at: Option<Timestamp>,
833}
834
835impl UploadFile {
836    async fn create(path: &Path, metadata: &Metadata) -> Result<()> {
837        let mut options = OpenOptions::from(file_options());
838        options.create_new(true).read(true).write(true);
839
840        let mut file = options.open(path).await.context(
841            ErrorKind::BackendFailure,
842            "creating local-fs resumable upload",
843        )?;
844        write_metadata_preamble(&mut file, metadata).await?;
845        file.sync_data().await.context(
846            ErrorKind::BackendFailure,
847            "syncing local-fs resumable upload",
848        )
849    }
850
851    async fn open(path: &Path) -> Result<Self> {
852        let file = match OpenOptions::new().read(true).write(true).open(path).await {
853            Ok(file) => file,
854            Err(error) if error.kind() == io::ErrorKind::NotFound => {
855                return Err(ErrorKind::UnknownUploadSession.into());
856            }
857            result => result.context(
858                ErrorKind::BackendFailure,
859                "opening local-fs resumable upload",
860            )?,
861        };
862        let mut reader = BufReader::new(file);
863        let (metadata, preamble_len) = read_metadata_preamble(&mut reader).await?;
864        let file = reader.into_inner();
865        let stored_size = file
866            .metadata()
867            .await
868            .context(
869                ErrorKind::BackendFailure,
870                "reading local-fs resumable upload size",
871            )?
872            .len();
873        let preamble_len = preamble_len as u64;
874        let payload_size = stored_size.checked_sub(preamble_len).ok_or_else(|| {
875            Error::new(
876                ErrorKind::CorruptData,
877                "reading truncated local-fs resumable upload",
878            )
879        })?;
880        Ok(Self {
881            file,
882            path: path.to_path_buf(),
883            stored_size,
884            expires_at: metadata.time_expires,
885            payload_size,
886        })
887    }
888
889    fn offset(&self) -> u64 {
890        self.payload_size
891    }
892
893    async fn append(&mut self, mut reader: impl AsyncRead + Unpin) -> Result<u64> {
894        self.file.seek(std::io::SeekFrom::End(0)).await.context(
895            ErrorKind::BackendFailure,
896            "seeking local-fs resumable upload",
897        )?;
898
899        let copied = match tokio::io::copy(&mut reader, &mut self.file).await {
900            Ok(copied) => copied,
901            Err(error) => {
902                let client_error = stream::unpack_client_error(&error);
903                self.file.sync_data().await.context(
904                    ErrorKind::BackendFailure,
905                    "syncing partial local-fs resumable chunk",
906                )?;
907                return Err(match client_error {
908                    Some(client_error) => Error::from(client_error),
909                    None => Error::with_context(
910                        ErrorKind::BackendFailure,
911                        "writing local-fs resumable chunk",
912                        error,
913                    ),
914                });
915            }
916        };
917
918        let payload_size = self.payload_size.checked_add(copied).ok_or_else(|| {
919            Error::new(
920                ErrorKind::BackendFailure,
921                "local-fs resumable payload size overflow",
922            )
923        })?;
924        let stored_size = self.stored_size.checked_add(copied).ok_or_else(|| {
925            Error::new(
926                ErrorKind::BackendFailure,
927                "local-fs resumable stored size overflow",
928            )
929        })?;
930        self.file.sync_data().await.context(
931            ErrorKind::BackendFailure,
932            "syncing local-fs resumable chunk",
933        )?;
934        self.payload_size = payload_size;
935        self.stored_size = stored_size;
936        Ok(payload_size)
937    }
938
939    /// Publishes the upload; requires a successful `append()` with no subsequent writes.
940    async fn publish(self, target: PathBuf) -> Result<(u64, Option<Timestamp>)> {
941        let Self {
942            file,
943            path,
944            stored_size,
945            expires_at,
946            ..
947        } = self;
948        drop(file);
949        tokio::fs::rename(path, target).await.context(
950            ErrorKind::BackendFailure,
951            "publishing local-fs resumable upload",
952        )?;
953        Ok((stored_size, expires_at))
954    }
955}
956
957struct ObjectFile {
958    metadata: Metadata,
959    preamble_len: u64,
960    payload_size: u64,
961    reader: BufReader<tokio::fs::File>,
962}
963
964async fn write_metadata_preamble<W>(writer: &mut W, metadata: &Metadata) -> Result<u64>
965where
966    W: tokio::io::AsyncWrite + Unpin,
967{
968    let metadata_json = serde_json::to_string(metadata)
969        .context(ErrorKind::Internal, "encoding local-fs object metadata")?;
970    writer.write_all(metadata_json.as_bytes()).await.context(
971        ErrorKind::BackendFailure,
972        "writing local-fs object metadata",
973    )?;
974    writer.write_all(b"\n").await.context(
975        ErrorKind::BackendFailure,
976        "writing local-fs object metadata",
977    )?;
978
979    Ok(metadata_json.len() as u64 + 1)
980}
981
982async fn read_metadata_preamble<R>(reader: &mut R) -> Result<(Metadata, usize)>
983where
984    R: tokio::io::AsyncBufRead + Unpin,
985{
986    let mut metadata_line = String::new();
987    let preamble_len = reader.read_line(&mut metadata_line).await.context(
988        ErrorKind::BackendFailure,
989        "reading local-fs object metadata",
990    )?;
991    let metadata = serde_json::from_str(metadata_line.trim_end())
992        .context(ErrorKind::CorruptData, "decoding local-fs object metadata")?;
993    Ok((metadata, preamble_len))
994}
995
996impl ObjectFile {
997    /// Opens an object file, returning `None` when it does not exist.
998    async fn try_open(path: &Path, access_time: Timestamp) -> Result<Option<Self>> {
999        let file = match OpenOptions::new().read(true).open(path).await {
1000            Ok(file) => file,
1001            Err(error) if error.kind() == io::ErrorKind::NotFound => return Ok(None),
1002            result => result.context(ErrorKind::BackendFailure, "opening local-fs object")?,
1003        };
1004
1005        let mut reader = BufReader::new(file);
1006        let (mut metadata, preamble_len) = read_metadata_preamble(&mut reader).await?;
1007
1008        if metadata.is_expired(access_time) {
1009            objectstore_log::debug!("Object found but past expiry");
1010            return Ok(None);
1011        }
1012
1013        let preamble_len = preamble_len as u64;
1014        let file_len = reader
1015            .get_ref()
1016            .metadata()
1017            .await
1018            .context(ErrorKind::BackendFailure, "reading local-fs object size")?
1019            .len();
1020
1021        let payload_size = file_len.checked_sub(preamble_len).ok_or_else(|| {
1022            Error::new(ErrorKind::CorruptData, "reading truncated local-fs object")
1023        })?;
1024
1025        metadata.size = Some(payload_size as usize);
1026
1027        Ok(Some(Self {
1028            metadata,
1029            preamble_len,
1030            payload_size,
1031            reader,
1032        }))
1033    }
1034}
1035
1036struct Draft {
1037    writer: BufWriter<tokio::fs::File>,
1038    path: tempfile::TempPath,
1039    target: PathBuf,
1040    /// Bytes the metadata preamble occupies, counted as it is written.
1041    preamble_len: u64,
1042}
1043
1044impl Draft {
1045    async fn create(target: &Path, metadata: &Metadata) -> Result<Self> {
1046        let parent = target.parent().unwrap().to_path_buf();
1047        let tempfile = tokio::task::spawn_blocking(move || create_tempfile(&parent))
1048            .await
1049            .context(
1050                ErrorKind::Internal,
1051                "waiting for local-fs object draft creation",
1052            )?
1053            .context(ErrorKind::BackendFailure, "creating local-fs object draft")?;
1054        let (file, path) = tempfile.into_parts();
1055
1056        let mut writer = BufWriter::new(tokio::fs::File::from_std(file));
1057        let preamble_len = write_metadata_preamble(&mut writer, metadata).await?;
1058
1059        Ok(Self {
1060            writer,
1061            path,
1062            target: target.to_path_buf(),
1063            preamble_len,
1064        })
1065    }
1066
1067    fn writer(&mut self) -> &mut BufWriter<tokio::fs::File> {
1068        &mut self.writer
1069    }
1070
1071    /// Bytes the metadata preamble occupies, for sizing the stored object.
1072    fn preamble_len(&self) -> u64 {
1073        self.preamble_len
1074    }
1075
1076    async fn prepare(&mut self) -> Result<()> {
1077        self.writer
1078            .flush()
1079            .await
1080            .context(ErrorKind::BackendFailure, "flushing local-fs object draft")?;
1081        self.writer
1082            .get_ref()
1083            .sync_data()
1084            .await
1085            .context(ErrorKind::BackendFailure, "syncing local-fs object draft")
1086    }
1087
1088    /// Publishes the draft; requires a successful `prepare()` with no subsequent writes.
1089    async fn publish(self) -> Result<()> {
1090        let Self {
1091            writer,
1092            path,
1093            target,
1094            ..
1095        } = self;
1096
1097        drop(writer);
1098        tokio::task::spawn_blocking(move || path.persist(target))
1099            .await
1100            .context(
1101                ErrorKind::Internal,
1102                "waiting to publish local-fs object draft",
1103            )?
1104            .context(
1105                ErrorKind::BackendFailure,
1106                "publishing local-fs object draft",
1107            )
1108    }
1109}
1110
1111fn create_tempfile(parent: &Path) -> io::Result<tempfile::NamedTempFile> {
1112    let mut builder = tempfile::Builder::new();
1113    builder.suffix(".draft");
1114
1115    builder.tempfile_in(parent)
1116}
1117
1118#[cfg(test)]
1119mod tests {
1120    use std::collections::HashMap;
1121    use std::num::NonZeroU32;
1122    use std::sync::Arc;
1123    use std::time::Duration;
1124
1125    use bytes::{Bytes, BytesMut};
1126    use futures_util::{TryStreamExt, stream as futures_stream};
1127    use objectstore_types::metadata::{Compression, ExpirationPolicy};
1128    use objectstore_types::scope::{Scope, Scopes};
1129    use objectstore_types::time::Timestamp;
1130
1131    #[cfg(feature = "storage-cogs")]
1132    use objectstore_inventory_tracker::OpType;
1133    #[cfg(feature = "storage-cogs")]
1134    use objectstore_inventory_tracker::test_utils::DummyProducer;
1135
1136    use super::*;
1137    use crate::id::ObjectContext;
1138    use crate::stream;
1139
1140    async fn upload_token(backend: &LocalFsBackend, id: &ObjectId, length: u64) -> BackendToken {
1141        backend
1142            .create_upload_session(id, &Metadata::default(), NonZeroU64::new(length).unwrap())
1143            .await
1144            .unwrap()
1145            .unwrap()
1146    }
1147
1148    #[tokio::test]
1149    async fn resumable_upload() {
1150        let (_tempdir, backend) = make_backend();
1151        let id = make_id();
1152        let metadata = Metadata {
1153            content_type: "text/resumable".into(),
1154            ..Default::default()
1155        };
1156
1157        // Create a session and verify its initial state on disk.
1158        let token = backend
1159            .create_upload_session(&id, &metadata, NonZeroU64::new(6).unwrap())
1160            .await
1161            .unwrap()
1162            .unwrap();
1163        let session = UploadSession::from_token(&token).unwrap();
1164        assert_eq!(session.total_length.get(), 6);
1165        let upload_path = backend.upload_path(session.upload_id);
1166        assert_eq!(
1167            upload_path.parent().unwrap().file_name().unwrap(),
1168            "uploads"
1169        );
1170        assert!(!backend.path(&id).exists());
1171
1172        // Upload the first chunk and verify that only the session advances.
1173        assert_eq!(
1174            backend.upload_offset(&id, &token).await.unwrap(),
1175            UploadProgress::Incomplete { offset: 0 }
1176        );
1177        assert_eq!(
1178            backend
1179                .put_chunk(&id, &token, 2, 0, stream::single(""))
1180                .await
1181                .unwrap(),
1182            UploadProgress::Incomplete { offset: 0 }
1183        );
1184        assert_eq!(
1185            backend
1186                .put_chunk(&id, &token, 0, 3, stream::single("abc"))
1187                .await
1188                .unwrap(),
1189            UploadProgress::Incomplete { offset: 3 }
1190        );
1191        assert!(!backend.path(&id).exists());
1192        assert_eq!(
1193            backend.upload_offset(&id, &token).await.unwrap(),
1194            UploadProgress::Incomplete { offset: 3 }
1195        );
1196
1197        // Zero-length chunks report current progress, while stale non-empty offsets fail.
1198        assert_eq!(
1199            backend
1200                .put_chunk(&id, &token, 0, 0, stream::single(""))
1201                .await
1202                .unwrap(),
1203            UploadProgress::Incomplete { offset: 3 }
1204        );
1205        let error = backend
1206            .put_chunk(&id, &token, 0, 3, stream::single("abc"))
1207            .await
1208            .unwrap_err();
1209        assert_eq!(error.kind(), ErrorKind::UploadOffsetMismatch { offset: 3 });
1210
1211        // Upload the final chunk and verify atomic publication removes the session.
1212        assert_eq!(
1213            backend
1214                .put_chunk(&id, &token, 3, 3, stream::single("def"))
1215                .await
1216                .unwrap(),
1217            UploadProgress::Complete
1218        );
1219        assert!(!upload_path.exists());
1220        let (stored_metadata, _, payload) = backend
1221            .get_object(&id, Timestamp::now(), None)
1222            .await
1223            .unwrap()
1224            .unwrap();
1225        assert_eq!(stored_metadata.content_type, metadata.content_type);
1226        assert_eq!(stream::read_to_vec(payload).await.unwrap(), b"abcdef");
1227        assert_eq!(
1228            backend.upload_offset(&id, &token).await.unwrap_err().kind(),
1229            ErrorKind::UnknownUploadSession
1230        );
1231    }
1232
1233    #[tokio::test]
1234    async fn resumable_publication_can_be_retried() {
1235        for query_offset in [false, true] {
1236            let (_tempdir, backend) = make_backend();
1237            let id = make_id();
1238            let token = upload_token(&backend, &id, 4).await;
1239            let object_path = backend.path(&id);
1240
1241            // A directory at the destination prevents rename after all bytes are persisted.
1242            tokio::fs::create_dir_all(&object_path).await.unwrap();
1243            let error = backend
1244                .put_chunk(&id, &token, 0, 4, stream::single("data"))
1245                .await
1246                .unwrap_err();
1247            assert_eq!(error.kind(), ErrorKind::BackendFailure);
1248            tokio::fs::remove_dir(&object_path).await.unwrap();
1249
1250            let progress = if query_offset {
1251                backend.upload_offset(&id, &token).await
1252            } else {
1253                backend
1254                    .put_chunk(&id, &token, 4, 0, stream::single(""))
1255                    .await
1256            };
1257            assert_eq!(progress.unwrap(), UploadProgress::Complete);
1258            let (_, _, payload) = backend
1259                .get_object(&id, Timestamp::now(), None)
1260                .await
1261                .unwrap()
1262                .unwrap();
1263            assert_eq!(stream::read_to_vec(payload).await.unwrap(), b"data");
1264            assert_eq!(
1265                backend.upload_offset(&id, &token).await.unwrap_err().kind(),
1266                ErrorKind::UnknownUploadSession
1267            );
1268        }
1269    }
1270
1271    #[tokio::test]
1272    async fn failed_resumable_chunk_preserves_partial_progress() {
1273        let (_tempdir, backend) = make_backend();
1274        let id = make_id();
1275        let token = upload_token(&backend, &id, 4).await;
1276
1277        // Persist a prefix, then disconnect after writing one byte of the next chunk.
1278        backend
1279            .put_chunk(&id, &token, 0, 2, stream::single("ab"))
1280            .await
1281            .unwrap();
1282        let broken = futures_stream::iter([
1283            Ok(Bytes::from_static(b"c")),
1284            Err(stream::ClientError::new(io::Error::other(
1285                "client disconnected",
1286            ))),
1287        ])
1288        .boxed();
1289        let error = backend
1290            .put_chunk(&id, &token, 2, 2, broken)
1291            .await
1292            .unwrap_err();
1293        assert_eq!(error.kind(), ErrorKind::ClientStream);
1294
1295        // Resume from the partial byte and verify the complete object.
1296        assert_eq!(
1297            backend.upload_offset(&id, &token).await.unwrap(),
1298            UploadProgress::Incomplete { offset: 3 }
1299        );
1300        assert_eq!(
1301            backend
1302                .put_chunk(&id, &token, 3, 1, stream::single("d"))
1303                .await
1304                .unwrap(),
1305            UploadProgress::Complete
1306        );
1307        let (_, _, payload) = backend
1308            .get_object(&id, Timestamp::now(), None)
1309            .await
1310            .unwrap()
1311            .unwrap();
1312        assert_eq!(stream::read_to_vec(payload).await.unwrap(), b"abcd");
1313    }
1314
1315    #[tokio::test]
1316    async fn concurrent_uploads_for_same_object_do_not_deadlock() {
1317        let (_tempdir, backend) = make_backend();
1318        let id = make_id();
1319        let first = upload_token(&backend, &id, 3).await;
1320        let second = upload_token(&backend, &id, 3).await;
1321
1322        // Complete both sessions concurrently through their shared object lock.
1323        let writes = async {
1324            tokio::join!(
1325                backend.put_chunk(&id, &first, 0, 3, stream::single("one")),
1326                backend.put_chunk(&id, &second, 0, 3, stream::single("two")),
1327            )
1328        };
1329        let (first_result, second_result) = tokio::time::timeout(Duration::from_secs(2), writes)
1330            .await
1331            .unwrap();
1332
1333        // Both requests finish without deadlock; the last publication wins.
1334        assert_eq!(first_result.unwrap(), UploadProgress::Complete);
1335        assert_eq!(second_result.unwrap(), UploadProgress::Complete);
1336        let (_, _, payload) = backend
1337            .get_object(&id, Timestamp::now(), None)
1338            .await
1339            .unwrap()
1340            .unwrap();
1341        let payload = stream::read_to_vec(payload).await.unwrap();
1342        assert!(payload == b"one" || payload == b"two");
1343    }
1344
1345    #[tokio::test]
1346    async fn stores_metadata() {
1347        let tempdir = tempfile::tempdir().unwrap();
1348        let backend = LocalFsBackend::new(
1349            FileSystemConfig {
1350                path: tempdir.path().to_path_buf(),
1351                cogs: None,
1352            },
1353            &ChangeStreamFactory::default(),
1354        );
1355
1356        let id = ObjectId::random(ObjectContext {
1357            usecase: "testing".into(),
1358            scopes: Scopes::from_iter([Scope::create("testing", "value").unwrap()]),
1359        });
1360
1361        let metadata = Metadata {
1362            content_type: "text/plain".into(),
1363            expiration_policy: ExpirationPolicy::TimeToIdle(Duration::from_hours(1)),
1364            time_created: Some(Timestamp::now()),
1365            time_expires: Some(Timestamp::now() + Duration::from_hours(1)),
1366            compression: Some(Compression::Zstd),
1367            origin: Some("203.0.113.42".into()),
1368            filename: Some("hello.txt".into()),
1369            custom: [("foo".into(), "bar".into())].into(),
1370            size: None,
1371        };
1372        backend
1373            .put_object(&id, &metadata, stream::single("oh hai!"), Timestamp::now())
1374            .await
1375            .unwrap();
1376
1377        let (read_metadata, _, stream) = backend
1378            .get_object(&id, Timestamp::now(), None)
1379            .await
1380            .unwrap()
1381            .unwrap();
1382        let file_contents: BytesMut = stream.try_collect().await.unwrap();
1383
1384        assert_eq!(
1385            read_metadata,
1386            Metadata {
1387                size: Some(file_contents.len()),
1388                ..metadata
1389            }
1390        );
1391        assert_eq!(file_contents.as_ref(), b"oh hai!");
1392
1393        let lock_path = backend.locks.path(&id);
1394        assert!(lock_path.exists());
1395        backend.delete_object(&id, Timestamp::now()).await.unwrap();
1396        assert!(lock_path.exists());
1397    }
1398
1399    #[tokio::test]
1400    async fn object_locks_coordinate_by_slot() {
1401        let tempdir = tempfile::tempdir().unwrap();
1402        let first_locks = ObjectLocks::new(tempdir.path());
1403        let second_locks = ObjectLocks::new(tempdir.path());
1404        let mut slots = HashMap::new();
1405        let (first_id, colliding_id) = (0..=65_536)
1406            .find_map(|key| {
1407                let id = ObjectId::from_parts("testing".into(), Scopes::empty(), key.to_string());
1408                let path = first_locks.path(&id);
1409                slots.insert(path, id.clone()).map(|first| (first, id))
1410            })
1411            .expect("65,537 keys must collide in 65,536 slots");
1412        let lock_path = first_locks.path(&first_id);
1413        let other_id = slots
1414            .values()
1415            .find(|id| first_locks.path(id) != lock_path)
1416            .unwrap();
1417        assert_eq!(second_locks.path(&colliding_id), lock_path);
1418
1419        let first_guard = first_locks.acquire(&first_id).await.unwrap();
1420        let mut waiter =
1421            tokio::spawn(async move { second_locks.acquire(&colliding_id).await.unwrap() });
1422        assert!(
1423            tokio::time::timeout(Duration::from_millis(50), &mut waiter)
1424                .await
1425                .is_err()
1426        );
1427
1428        let other_guard =
1429            tokio::time::timeout(Duration::from_secs(1), first_locks.acquire(other_id))
1430                .await
1431                .unwrap()
1432                .unwrap();
1433        drop(other_guard);
1434        drop(first_guard);
1435        let guard = tokio::time::timeout(Duration::from_secs(1), waiter)
1436            .await
1437            .unwrap()
1438            .unwrap();
1439        drop(guard);
1440        assert!(lock_path.exists());
1441    }
1442
1443    #[tokio::test]
1444    async fn missing_object_expiry_does_not_block_descendant() {
1445        let (_tempdir, backend) = make_backend();
1446        let id = ObjectId::from_parts("testing".into(), Scopes::empty(), "foo".into());
1447        assert!(
1448            !backend
1449                .set_expiry(
1450                    &id,
1451                    Timestamp::now() + Duration::from_hours(1),
1452                    Timestamp::now()
1453                )
1454                .await
1455                .unwrap()
1456        );
1457        let descendant = ObjectId::new(id.context.clone(), "foo/bar".into());
1458        backend
1459            .put_object(
1460                &descendant,
1461                &Metadata::default(),
1462                stream::single("payload"),
1463                Timestamp::now(),
1464            )
1465            .await
1466            .unwrap();
1467        let (_, _, payload) = backend
1468            .get_object(&descendant, Timestamp::now(), None)
1469            .await
1470            .unwrap()
1471            .unwrap();
1472        assert_eq!(stream::read_to_vec(payload).await.unwrap(), b"payload");
1473    }
1474
1475    #[tokio::test]
1476    async fn failed_put_preserves_published_object() {
1477        let (_tempdir, backend) = make_backend();
1478        let id = make_id();
1479        let original_metadata = Metadata {
1480            content_type: "text/original".into(),
1481            ..Default::default()
1482        };
1483        backend
1484            .put_object(
1485                &id,
1486                &original_metadata,
1487                stream::single("original"),
1488                Timestamp::now(),
1489            )
1490            .await
1491            .unwrap();
1492
1493        let replacement_metadata = Metadata {
1494            content_type: "text/replacement".into(),
1495            ..Default::default()
1496        };
1497        let replacement = futures_stream::iter([
1498            Ok(Bytes::from_static(b"partial")),
1499            Err(stream::ClientError::new(io::Error::other(
1500                "replacement stream failed",
1501            ))),
1502        ])
1503        .boxed();
1504        let error = backend
1505            .put_object(&id, &replacement_metadata, replacement, Timestamp::now())
1506            .await
1507            .unwrap_err();
1508        assert_eq!(error.kind(), ErrorKind::ClientStream);
1509
1510        let (metadata, _, payload) = backend
1511            .get_object(&id, Timestamp::now(), None)
1512            .await
1513            .unwrap()
1514            .unwrap();
1515        assert_eq!(metadata.content_type, original_metadata.content_type);
1516        assert_eq!(stream::read_to_vec(payload).await.unwrap(), b"original");
1517
1518        assert_eq!(draft_count(&backend, &id), 0);
1519    }
1520
1521    #[tokio::test]
1522    async fn cancelled_put_removes_draft() {
1523        let (_tempdir, backend) = make_backend();
1524        let backend = Arc::new(backend);
1525        let id = make_id();
1526        backend
1527            .put_object(
1528                &id,
1529                &Metadata::default(),
1530                stream::single("original"),
1531                Timestamp::now(),
1532            )
1533            .await
1534            .unwrap();
1535
1536        let (writing, writing_started) = tokio::sync::oneshot::channel();
1537        let replacement = futures_stream::once(async move {
1538            let _ = writing.send(());
1539            Ok(Bytes::from_static(b"partial"))
1540        })
1541        .chain(futures_stream::pending())
1542        .boxed();
1543
1544        let task = tokio::spawn({
1545            let backend = Arc::clone(&backend);
1546            let id = id.clone();
1547            async move {
1548                backend
1549                    .put_object(&id, &Metadata::default(), replacement, Timestamp::now())
1550                    .await
1551            }
1552        });
1553
1554        writing_started.await.unwrap();
1555        assert_eq!(draft_count(&backend, &id), 1);
1556        task.abort();
1557        assert!(task.await.unwrap_err().is_cancelled());
1558
1559        let (_, _, payload) = backend
1560            .get_object(&id, Timestamp::now(), None)
1561            .await
1562            .unwrap()
1563            .unwrap();
1564        assert_eq!(stream::read_to_vec(payload).await.unwrap(), b"original");
1565
1566        assert_eq!(draft_count(&backend, &id), 0);
1567    }
1568
1569    #[tokio::test]
1570    async fn set_expiry() {
1571        let (_tempdir, backend) = make_backend();
1572        let id = make_id();
1573        let old_expiry = Timestamp::now() + Duration::from_hours(1);
1574        let metadata = Metadata {
1575            expiration_policy: ExpirationPolicy::TimeToIdle(Duration::from_hours(1)),
1576            time_expires: Some(old_expiry),
1577            custom: [("preserved".into(), "yes".into())].into(),
1578            ..Default::default()
1579        };
1580        backend
1581            .put_object(&id, &metadata, stream::single("payload"), Timestamp::now())
1582            .await
1583            .unwrap();
1584
1585        let requested = old_expiry + Duration::from_hours(1) + Duration::from_nanos(999);
1586        assert!(
1587            backend
1588                .set_expiry(&id, requested, Timestamp::now())
1589                .await
1590                .unwrap()
1591        );
1592        let (updated, _, payload) = backend
1593            .get_object(&id, Timestamp::now(), None)
1594            .await
1595            .unwrap()
1596            .unwrap();
1597        assert_eq!(updated.expiration_policy, metadata.expiration_policy);
1598        assert_eq!(updated.custom, metadata.custom);
1599        assert_eq!(updated.time_expires, Some(requested));
1600        assert_eq!(stream::read_to_vec(payload).await.unwrap(), b"payload");
1601
1602        assert_eq!(draft_count(&backend, &id), 0);
1603    }
1604
1605    #[tokio::test]
1606    async fn expired_object() {
1607        let (_tempdir, backend) = make_backend();
1608        let id = make_id();
1609        let metadata = Metadata {
1610            expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_hours(1)),
1611            time_expires: Some(Timestamp::now() - Duration::from_secs(1)),
1612            ..Default::default()
1613        };
1614        backend
1615            .put_object(&id, &metadata, stream::single("expired"), Timestamp::now())
1616            .await
1617            .unwrap();
1618
1619        assert!(
1620            backend
1621                .get_object(&id, Timestamp::now(), None)
1622                .await
1623                .unwrap()
1624                .is_none()
1625        );
1626        assert!(
1627            !backend
1628                .set_expiry(
1629                    &id,
1630                    Timestamp::now() + Duration::from_hours(1),
1631                    Timestamp::now()
1632                )
1633                .await
1634                .unwrap()
1635        );
1636    }
1637
1638    #[tokio::test]
1639    async fn get_metadata_returns_metadata() {
1640        let tempdir = tempfile::tempdir().unwrap();
1641        let backend = LocalFsBackend::new(
1642            FileSystemConfig {
1643                path: tempdir.path().to_path_buf(),
1644                cogs: None,
1645            },
1646            &ChangeStreamFactory::default(),
1647        );
1648
1649        let id = ObjectId::random(ObjectContext {
1650            usecase: "testing".into(),
1651            scopes: Scopes::from_iter([Scope::create("testing", "value").unwrap()]),
1652        });
1653
1654        let metadata = Metadata {
1655            content_type: "text/plain".into(),
1656            compression: Some(Compression::Zstd),
1657            origin: Some("203.0.113.42".into()),
1658            custom: [("foo".into(), "bar".into())].into(),
1659            ..Default::default()
1660        };
1661        backend
1662            .put_object(&id, &metadata, stream::single("oh hai!"), Timestamp::now())
1663            .await
1664            .unwrap();
1665
1666        let read_metadata = backend
1667            .get_metadata(&id, Timestamp::now())
1668            .await
1669            .unwrap()
1670            .unwrap();
1671        assert_eq!(
1672            read_metadata,
1673            Metadata {
1674                size: Some(7),
1675                ..metadata
1676            }
1677        );
1678    }
1679
1680    #[tokio::test]
1681    async fn get_metadata_nonexistent() {
1682        let tempdir = tempfile::tempdir().unwrap();
1683        let backend = LocalFsBackend::new(
1684            FileSystemConfig {
1685                path: tempdir.path().to_path_buf(),
1686                cogs: None,
1687            },
1688            &ChangeStreamFactory::default(),
1689        );
1690
1691        let id = ObjectId::random(ObjectContext {
1692            usecase: "testing".into(),
1693            scopes: Scopes::from_iter([Scope::create("testing", "value").unwrap()]),
1694        });
1695
1696        let result = backend.get_metadata(&id, Timestamp::now()).await.unwrap();
1697        assert!(result.is_none());
1698    }
1699
1700    fn make_id() -> ObjectId {
1701        ObjectId::random(ObjectContext {
1702            usecase: "testing".into(),
1703            scopes: Scopes::from_iter([Scope::create("testing", "value").unwrap()]),
1704        })
1705    }
1706
1707    #[cfg(feature = "storage-cogs")]
1708    fn make_backend_with_change_stream() -> (tempfile::TempDir, LocalFsBackend, DummyProducer) {
1709        let tempdir = tempfile::tempdir().unwrap();
1710        let (streams, producer) = crate::change_stream::dummy_factory();
1711        let backend = LocalFsBackend::new(
1712            FileSystemConfig {
1713                path: tempdir.path().to_path_buf(),
1714                cogs: Some(CostTrackerStreamConfig {
1715                    shared_resource_id: "filesystem_objectstore".into(),
1716                    sample_rate: 1.0,
1717                }),
1718            },
1719            &streams,
1720        );
1721        (tempdir, backend, producer)
1722    }
1723
1724    fn make_backend() -> (tempfile::TempDir, LocalFsBackend) {
1725        let tempdir = tempfile::tempdir().unwrap();
1726        let backend = LocalFsBackend::new(
1727            FileSystemConfig {
1728                path: tempdir.path().to_path_buf(),
1729                cogs: None,
1730            },
1731            &ChangeStreamFactory::default(),
1732        );
1733        (tempdir, backend)
1734    }
1735
1736    fn draft_count(backend: &LocalFsBackend, id: &ObjectId) -> usize {
1737        let object_path = backend.path(id);
1738        std::fs::read_dir(object_path.parent().unwrap())
1739            .unwrap()
1740            .flat_map(|entry| entry.map(|entry| entry.path()))
1741            .filter(|path| path.to_string_lossy().ends_with(".draft"))
1742            .count()
1743    }
1744
1745    #[tokio::test]
1746    async fn multipart_single_part() {
1747        let (_tempdir, backend) = make_backend();
1748        let id = make_id();
1749        let metadata = Metadata {
1750            content_type: "text/plain".into(),
1751            expiration_policy: ExpirationPolicy::TimeToIdle(Duration::from_hours(1)),
1752            origin: Some("203.0.113.42".into()),
1753            custom: [("foo".into(), "bar".into())].into(),
1754            ..Default::default()
1755        };
1756
1757        let upload_id = backend.initiate_multipart(&id, &metadata).await.unwrap();
1758
1759        let data = b"hello, multipart world!";
1760        let etag = backend
1761            .upload_part(
1762                &id,
1763                &upload_id,
1764                NonZeroU32::new(1).unwrap(),
1765                data.len() as u64,
1766                None,
1767                stream::single(data.to_vec()),
1768            )
1769            .await
1770            .unwrap();
1771
1772        let result = backend
1773            .complete_multipart(
1774                &id,
1775                &upload_id,
1776                vec![CompletedPart {
1777                    part_number: NonZeroU32::new(1).unwrap(),
1778                    etag,
1779                }],
1780                Timestamp::now(),
1781            )
1782            .await
1783            .unwrap();
1784        assert!(result.is_none(), "expected no error on complete");
1785
1786        let (meta, _, body) = backend
1787            .get_object(&id, Timestamp::now(), None)
1788            .await
1789            .unwrap()
1790            .unwrap();
1791        let payload: BytesMut = body.try_collect().await.unwrap();
1792        assert_eq!(payload.as_ref(), data);
1793        assert_eq!(meta.content_type, "text/plain".to_string());
1794        assert_eq!(
1795            meta.expiration_policy,
1796            ExpirationPolicy::TimeToIdle(Duration::from_hours(1))
1797        );
1798        assert_eq!(meta.origin, Some("203.0.113.42".into()));
1799        assert_eq!(meta.custom, [("foo".into(), "bar".into())].into());
1800    }
1801
1802    #[tokio::test]
1803    async fn multipart_multiple_parts() {
1804        let (_tempdir, backend) = make_backend();
1805        let id = make_id();
1806        let metadata = Metadata::default();
1807
1808        let upload_id = backend.initiate_multipart(&id, &metadata).await.unwrap();
1809
1810        let part1 = b"aaaa".to_vec();
1811        let part2 = b"bbbb".to_vec();
1812        let part3 = b"cc".to_vec();
1813
1814        let etag1 = backend
1815            .upload_part(
1816                &id,
1817                &upload_id,
1818                NonZeroU32::new(1).unwrap(),
1819                part1.len() as u64,
1820                None,
1821                stream::single(part1.clone()),
1822            )
1823            .await
1824            .unwrap();
1825        let etag2 = backend
1826            .upload_part(
1827                &id,
1828                &upload_id,
1829                NonZeroU32::new(2).unwrap(),
1830                part2.len() as u64,
1831                None,
1832                stream::single(part2.clone()),
1833            )
1834            .await
1835            .unwrap();
1836        let etag3 = backend
1837            .upload_part(
1838                &id,
1839                &upload_id,
1840                NonZeroU32::new(3).unwrap(),
1841                part3.len() as u64,
1842                None,
1843                stream::single(part3.clone()),
1844            )
1845            .await
1846            .unwrap();
1847
1848        let result = backend
1849            .complete_multipart(
1850                &id,
1851                &upload_id,
1852                vec![
1853                    CompletedPart {
1854                        part_number: NonZeroU32::new(1).unwrap(),
1855                        etag: etag1,
1856                    },
1857                    CompletedPart {
1858                        part_number: NonZeroU32::new(2).unwrap(),
1859                        etag: etag2,
1860                    },
1861                    CompletedPart {
1862                        part_number: NonZeroU32::new(3).unwrap(),
1863                        etag: etag3,
1864                    },
1865                ],
1866                Timestamp::now(),
1867            )
1868            .await
1869            .unwrap();
1870        assert!(result.is_none());
1871
1872        let (_, _, body) = backend
1873            .get_object(&id, Timestamp::now(), None)
1874            .await
1875            .unwrap()
1876            .unwrap();
1877        let payload: BytesMut = body.try_collect().await.unwrap();
1878        assert_eq!(payload.as_ref(), b"aaaabbbbcc");
1879    }
1880
1881    #[tokio::test]
1882    async fn multipart_list_parts() {
1883        let (_tempdir, backend) = make_backend();
1884        let id = make_id();
1885        let metadata = Metadata::default();
1886
1887        let upload_id = backend.initiate_multipart(&id, &metadata).await.unwrap();
1888
1889        let etag1 = backend
1890            .upload_part(
1891                &id,
1892                &upload_id,
1893                NonZeroU32::new(1).unwrap(),
1894                3,
1895                None,
1896                stream::single(b"aaa".to_vec()),
1897            )
1898            .await
1899            .unwrap();
1900        let etag2 = backend
1901            .upload_part(
1902                &id,
1903                &upload_id,
1904                NonZeroU32::new(2).unwrap(),
1905                3,
1906                None,
1907                stream::single(b"bbb".to_vec()),
1908            )
1909            .await
1910            .unwrap();
1911
1912        let list = backend
1913            .list_parts(&id, &upload_id, None, None)
1914            .await
1915            .unwrap();
1916        assert_eq!(list.parts.len(), 2);
1917        assert_eq!(list.parts[0].part_number.get(), 1);
1918        assert_eq!(list.parts[0].etag, etag1);
1919        assert_eq!(list.parts[0].size, 3);
1920        assert_eq!(list.parts[1].part_number.get(), 2);
1921        assert_eq!(list.parts[1].etag, etag2);
1922        assert_eq!(list.parts[1].size, 3);
1923
1924        // Pagination
1925        let page1 = backend
1926            .list_parts(&id, &upload_id, Some(1), None)
1927            .await
1928            .unwrap();
1929        assert_eq!(page1.parts.len(), 1);
1930        assert_eq!(page1.parts[0].part_number.get(), 1);
1931        assert!(page1.is_truncated);
1932        assert!(page1.next_part_number_marker.is_some());
1933
1934        let page2 = backend
1935            .list_parts(&id, &upload_id, Some(1), page1.next_part_number_marker)
1936            .await
1937            .unwrap();
1938        assert_eq!(page2.parts.len(), 1);
1939        assert_eq!(page2.parts[0].part_number.get(), 2);
1940
1941        backend.abort_multipart(&id, &upload_id).await.unwrap();
1942    }
1943
1944    #[tokio::test]
1945    async fn get_object_range_bounded() {
1946        let (_tempdir, backend) = make_backend();
1947        let id = make_id();
1948        let metadata = Metadata::default();
1949
1950        let payload = b"Hello, range requests!";
1951        backend
1952            .put_object(
1953                &id,
1954                &metadata,
1955                stream::single(payload.to_vec()),
1956                Timestamp::now(),
1957            )
1958            .await
1959            .unwrap();
1960
1961        // Request bytes 7-11 → "range"
1962        let (_, content_range, body) = backend
1963            .get_object(&id, Timestamp::now(), Some(ByteRange::Bounded(7, 11)))
1964            .await
1965            .unwrap()
1966            .unwrap();
1967        let data: BytesMut = body.try_collect().await.unwrap();
1968
1969        assert_eq!(data.as_ref(), b"range");
1970        let content_range = content_range.unwrap();
1971        assert_eq!(content_range.start, 7);
1972        assert_eq!(content_range.end, 11);
1973        assert_eq!(content_range.total, payload.len() as u64);
1974    }
1975
1976    #[tokio::test]
1977    async fn get_object_range_from() {
1978        let (_tempdir, backend) = make_backend();
1979        let id = make_id();
1980        let metadata = Metadata::default();
1981
1982        let payload = b"Hello, range requests!";
1983        backend
1984            .put_object(
1985                &id,
1986                &metadata,
1987                stream::single(payload.to_vec()),
1988                Timestamp::now(),
1989            )
1990            .await
1991            .unwrap();
1992
1993        // Request bytes 7- → "range requests!"
1994        let (_, content_range, body) = backend
1995            .get_object(&id, Timestamp::now(), Some(ByteRange::From(7)))
1996            .await
1997            .unwrap()
1998            .unwrap();
1999        let data: BytesMut = body.try_collect().await.unwrap();
2000
2001        assert_eq!(data.as_ref(), b"range requests!");
2002        let content_range = content_range.unwrap();
2003        assert_eq!(content_range.start, 7);
2004        assert_eq!(content_range.end, 21);
2005        assert_eq!(content_range.total, payload.len() as u64);
2006    }
2007
2008    #[tokio::test]
2009    async fn get_object_range_last() {
2010        let (_tempdir, backend) = make_backend();
2011        let id = make_id();
2012        let metadata = Metadata::default();
2013
2014        let payload = b"Hello, range requests!";
2015        backend
2016            .put_object(
2017                &id,
2018                &metadata,
2019                stream::single(payload.to_vec()),
2020                Timestamp::now(),
2021            )
2022            .await
2023            .unwrap();
2024
2025        // Request last 9 bytes → "requests!"
2026        let (_, content_range, body) = backend
2027            .get_object(&id, Timestamp::now(), Some(ByteRange::Last(9)))
2028            .await
2029            .unwrap()
2030            .unwrap();
2031        let data: BytesMut = body.try_collect().await.unwrap();
2032
2033        assert_eq!(data.as_ref(), b"requests!");
2034        let content_range = content_range.unwrap();
2035        assert_eq!(content_range.start, 13);
2036        assert_eq!(content_range.end, 21);
2037        assert_eq!(content_range.total, payload.len() as u64);
2038    }
2039
2040    #[tokio::test]
2041    async fn get_object_range_unsatisfiable() {
2042        let (_tempdir, backend) = make_backend();
2043        let id = make_id();
2044        let metadata = Metadata::default();
2045
2046        backend
2047            .put_object(
2048                &id,
2049                &metadata,
2050                stream::single(b"short".to_vec()),
2051                Timestamp::now(),
2052            )
2053            .await
2054            .unwrap();
2055
2056        match backend
2057            .get_object(&id, Timestamp::now(), Some(ByteRange::From(100)))
2058            .await
2059        {
2060            Err(error) if matches!(error.kind(), ErrorKind::RangeNotSatisfiable { total: 5 }) => {}
2061            Err(other) => panic!("expected RangeNotSatisfiable, got: {other:?}"),
2062            Ok(_) => panic!("expected RangeNotSatisfiable, got Ok"),
2063        }
2064    }
2065
2066    #[tokio::test]
2067    async fn multipart_abort() {
2068        let (_tempdir, backend) = make_backend();
2069        let id = make_id();
2070        let metadata = Metadata::default();
2071
2072        let upload_id = backend.initiate_multipart(&id, &metadata).await.unwrap();
2073
2074        backend
2075            .upload_part(
2076                &id,
2077                &upload_id,
2078                NonZeroU32::new(1).unwrap(),
2079                5,
2080                None,
2081                stream::single(b"hello".to_vec()),
2082            )
2083            .await
2084            .unwrap();
2085
2086        backend.abort_multipart(&id, &upload_id).await.unwrap();
2087
2088        let result = backend
2089            .get_object(&id, Timestamp::now(), None)
2090            .await
2091            .unwrap();
2092        assert!(result.is_none(), "object should not exist after abort");
2093    }
2094
2095    #[tokio::test]
2096    async fn multipart_invalid_etag() {
2097        let (_tempdir, backend) = make_backend();
2098        let id = make_id();
2099        let metadata = Metadata::default();
2100
2101        let upload_id = backend.initiate_multipart(&id, &metadata).await.unwrap();
2102
2103        let etag = backend
2104            .upload_part(
2105                &id,
2106                &upload_id,
2107                NonZeroU32::new(1).unwrap(),
2108                5,
2109                None,
2110                stream::single(b"hello".to_vec()),
2111            )
2112            .await
2113            .unwrap();
2114
2115        let result = backend
2116            .complete_multipart(
2117                &id,
2118                &upload_id,
2119                vec![CompletedPart {
2120                    part_number: NonZeroU32::new(1).unwrap(),
2121                    etag: "wrong-etag".into(),
2122                }],
2123                Timestamp::now(),
2124            )
2125            .await
2126            .unwrap();
2127        assert!(result.is_some(), "expected error for bad etag");
2128        assert_eq!(result.unwrap().code, "InvalidPart");
2129
2130        // Upload must survive a failed complete so the client can retry.
2131        let result = backend
2132            .complete_multipart(
2133                &id,
2134                &upload_id,
2135                vec![CompletedPart {
2136                    part_number: NonZeroU32::new(1).unwrap(),
2137                    etag,
2138                }],
2139                Timestamp::now(),
2140            )
2141            .await
2142            .unwrap();
2143        assert!(result.is_none(), "retry with correct etag should succeed");
2144    }
2145
2146    #[tokio::test]
2147    async fn multipart_missing_part() {
2148        let (_tempdir, backend) = make_backend();
2149        let id = make_id();
2150        let metadata = Metadata::default();
2151
2152        let upload_id = backend.initiate_multipart(&id, &metadata).await.unwrap();
2153
2154        let etag = backend
2155            .upload_part(
2156                &id,
2157                &upload_id,
2158                NonZeroU32::new(1).unwrap(),
2159                5,
2160                None,
2161                stream::single(b"hello".to_vec()),
2162            )
2163            .await
2164            .unwrap();
2165
2166        let result = backend
2167            .complete_multipart(
2168                &id,
2169                &upload_id,
2170                vec![CompletedPart {
2171                    part_number: NonZeroU32::new(99).unwrap(),
2172                    etag: "whatever".into(),
2173                }],
2174                Timestamp::now(),
2175            )
2176            .await
2177            .unwrap();
2178        assert!(result.is_some(), "expected error for missing part");
2179        assert_eq!(result.unwrap().code, "InvalidPart");
2180
2181        // Upload must survive a failed complete so the client can retry.
2182        let result = backend
2183            .complete_multipart(
2184                &id,
2185                &upload_id,
2186                vec![CompletedPart {
2187                    part_number: NonZeroU32::new(1).unwrap(),
2188                    etag,
2189                }],
2190                Timestamp::now(),
2191            )
2192            .await
2193            .unwrap();
2194        assert!(result.is_none(), "retry with correct part should succeed");
2195    }
2196
2197    #[cfg(feature = "storage-cogs")]
2198    #[tokio::test]
2199    async fn change_stream_reports_size_written_to_disk() {
2200        let (tempdir, backend, producer) = make_backend_with_change_stream();
2201        let id = make_id();
2202        let metadata = Metadata {
2203            expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_hours(1)),
2204            time_expires: Some(Timestamp::now() + Duration::from_hours(1)),
2205            ..Default::default()
2206        };
2207        let payload = b"oh hai!";
2208
2209        backend
2210            .put_object(
2211                &id,
2212                &metadata,
2213                stream::single(payload.to_vec()),
2214                Timestamp::now(),
2215            )
2216            .await
2217            .unwrap();
2218
2219        let file = tokio::fs::read(tempdir.path().join(id.as_storage_path().to_string()))
2220            .await
2221            .unwrap();
2222
2223        let records = producer.records();
2224        assert_eq!(records.len(), 1);
2225        assert_eq!(records[0].op_type, OpType::Write);
2226        assert_eq!(records[0].shared_resource_id, "filesystem_objectstore");
2227        assert_eq!(records[0].app_feature, "testing");
2228        assert_eq!(
2229            records[0].size,
2230            Some(file.len() as u64),
2231            "the metadata header line counts towards the reported size"
2232        );
2233        assert!(records[0].expiration_time.is_some());
2234    }
2235
2236    #[cfg(feature = "storage-cogs")]
2237    #[tokio::test]
2238    async fn change_stream_reports_resumable_upload_size() {
2239        let (tempdir, backend, producer) = make_backend_with_change_stream();
2240        let id = make_id();
2241        let metadata = Metadata {
2242            expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_hours(1)),
2243            time_expires: Some(Timestamp::now() + Duration::from_hours(1)),
2244            ..Default::default()
2245        };
2246        let payload = b"oh hai!";
2247        let token = backend
2248            .create_upload_session(
2249                &id,
2250                &metadata,
2251                NonZeroU64::new(payload.len() as u64).unwrap(),
2252            )
2253            .await
2254            .unwrap()
2255            .unwrap();
2256
2257        // Session creation alone does not report a stored object.
2258        assert!(producer.records().is_empty());
2259
2260        // Completing the upload reports the published file's full stored size.
2261        assert_eq!(
2262            backend
2263                .put_chunk(
2264                    &id,
2265                    &token,
2266                    0,
2267                    payload.len() as u64,
2268                    stream::single(payload.to_vec()),
2269                )
2270                .await
2271                .unwrap(),
2272            UploadProgress::Complete
2273        );
2274
2275        let file = tokio::fs::read(tempdir.path().join(id.as_storage_path().to_string()))
2276            .await
2277            .unwrap();
2278        let records = producer.records();
2279        assert_eq!(records.len(), 1);
2280        assert_eq!(records[0].op_type, OpType::Write);
2281        assert_eq!(records[0].size, Some(file.len() as u64));
2282        assert!(records[0].expiration_time.is_some());
2283    }
2284
2285    #[cfg(feature = "storage-cogs")]
2286    #[tokio::test]
2287    async fn change_stream_reports_assembled_multipart_size() {
2288        let (tempdir, backend, producer) = make_backend_with_change_stream();
2289        let id = make_id();
2290        let metadata = Metadata::default();
2291
2292        let upload_id = backend.initiate_multipart(&id, &metadata).await.unwrap();
2293
2294        let mut completed = Vec::new();
2295        for (number, payload) in [(1u32, b"aaaa".to_vec()), (2, b"bbbb".to_vec())] {
2296            let part_number = NonZeroU32::new(number).unwrap();
2297            let etag = backend
2298                .upload_part(
2299                    &id,
2300                    &upload_id,
2301                    part_number,
2302                    payload.len() as u64,
2303                    None,
2304                    stream::single(payload),
2305                )
2306                .await
2307                .unwrap();
2308            completed.push(CompletedPart { part_number, etag });
2309        }
2310
2311        // Uploading parts is intermediate state, so nothing is reported until completion.
2312        assert!(producer.records().is_empty());
2313
2314        assert!(
2315            backend
2316                .complete_multipart(&id, &upload_id, completed, Timestamp::now())
2317                .await
2318                .unwrap()
2319                .is_none()
2320        );
2321
2322        let file = tokio::fs::read(tempdir.path().join(id.as_storage_path().to_string()))
2323            .await
2324            .unwrap();
2325
2326        let records = producer.records();
2327        assert_eq!(records.len(), 1);
2328        assert_eq!(records[0].op_type, OpType::Write);
2329        assert_eq!(
2330            records[0].size,
2331            Some(file.len() as u64),
2332            "the assembled object counts its metadata header and every part payload"
2333        );
2334    }
2335
2336    #[cfg(feature = "storage-cogs")]
2337    #[tokio::test]
2338    async fn change_stream_reports_deletes_on_success() {
2339        let (_tempdir, backend, producer) = make_backend_with_change_stream();
2340        let id = make_id();
2341
2342        // Try to delete a non-existent object. Don't emit a message.
2343        backend
2344            .delete_object(&id, Timestamp::now())
2345            .await
2346            .expect("deleting a non-existent object returns Ok(())");
2347        assert!(producer.records().is_empty());
2348
2349        backend
2350            .put_object(
2351                &id,
2352                &Metadata::default(),
2353                stream::single(b"hi".to_vec()),
2354                Timestamp::now(),
2355            )
2356            .await
2357            .unwrap();
2358        producer.clear();
2359
2360        backend.delete_object(&id, Timestamp::now()).await.unwrap();
2361
2362        let records = producer.records();
2363        assert_eq!(records.len(), 1);
2364        assert_eq!(records[0].op_type, OpType::Delete);
2365    }
2366}