1use 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
55fn 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
66async 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#[derive(Debug, Clone, serde::Deserialize, serde::Serialize)]
88pub struct FileSystemConfig {
89 pub path: PathBuf,
103
104 #[serde(default, skip_serializing_if = "Option::is_none")]
115 pub cogs: Option<CostTrackerStreamConfig>,
116}
117
118#[derive(Debug)]
120pub struct LocalFsBackend {
121 path: PathBuf,
122 locks: ObjectLocks,
123
124 change_stream: Arc<dyn ChangeStream>,
125}
126
127impl LocalFsBackend {
128 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 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); }
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 #[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 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 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 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 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 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
744const 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 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 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
824struct UploadFile {
826 file: tokio::fs::File,
827 path: PathBuf,
828 payload_size: u64,
830 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 assert!(producer.records().is_empty());
2259
2260 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 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 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}