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