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 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
56fn 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
67async 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#[derive(Debug, Clone, serde::Deserialize, serde::Serialize)]
89pub struct FileSystemConfig {
90 pub path: PathBuf,
104
105 #[serde(default, skip_serializing_if = "Option::is_none")]
116 pub cogs: Option<CostTrackerStreamConfig>,
117}
118
119#[derive(Debug)]
121pub struct LocalFsBackend {
122 path: PathBuf,
123 locks: ObjectLocks,
124
125 change_stream: Arc<dyn ChangeStream>,
126}
127
128impl LocalFsBackend {
129 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 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)); }
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 #[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 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 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 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 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 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
756const 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 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 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
812struct UploadFile {
814 file: tokio::fs::File,
815 path: PathBuf,
816 payload_size: u64,
818 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 assert!(producer.records().is_empty());
2280
2281 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 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 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}