diff --git a/objectstore-service/docs/architecture.md b/objectstore-service/docs/architecture.md index f195610f..145ab78c 100644 --- a/objectstore-service/docs/architecture.md +++ b/objectstore-service/docs/architecture.md @@ -182,8 +182,8 @@ See also: [`objectstore_inventory_tracker`] documentation. # Change Streams Every backend publishes the changes it makes to the objects it stores as a -[`ChangeStream`](change_stream::ChangeStream). It is a fire-and-forget, -per-backend feed of three operations: +[`ChangeStream`](change_stream::ChangeStream). It is a per-backend feed of +three operations: - `write(id, size, expires_at)`: `id` now occupies `size` bytes. Used for both new objects and overwrites. @@ -191,6 +191,13 @@ per-backend feed of three operations: unchanged. In practice this is a TTI bump. - `delete(id)`: `id` was deleted explicitly. +Writes and updates are reported twice: `begin_*` before the storage operation +and `commit_*` after it succeeds. Deletes are only committed. A failed +`begin_*` fails the request before storage is touched, while `commit_*` cannot +fail. Resumable and multipart uploads call `begin_write` when they are created. +A `begin_*` is not always followed by a `commit_*`, so `begin_*` gives an upper +bound on what is stored, e.g. for cleaning up expired data. + The stream describes physical storage per backend. When using [`TieredStorage`](backend::tiered::TieredStorage), objects that are stored in long-term storage will emit a change record for the actual object in long-term diff --git a/objectstore-service/src/backend/bigtable.rs b/objectstore-service/src/backend/bigtable.rs index 04df7e5c..392424b6 100644 --- a/objectstore-service/src/backend/bigtable.rs +++ b/objectstore-service/src/backend/bigtable.rs @@ -999,6 +999,9 @@ impl Backend for BigTableBackend { _access_time: Timestamp, ) -> Result { objectstore_log::debug!("Writing to Bigtable backend"); + self.change_stream + .begin_write(id, metadata.time_expires) + .await?; let path = id.as_storage_path().to_string().into_bytes(); let mut payload = ChunkedBytes::new(0); @@ -1009,7 +1012,9 @@ impl Backend for BigTableBackend { let (_, size) = self .put_row(path, metadata.clone(), payload.into_bytes().into(), "put") .await?; - self.change_stream.write(id, size, metadata.time_expires); + self.change_stream + .commit_write(id, size, metadata.time_expires) + .await; Ok(()) } @@ -1063,7 +1068,7 @@ impl Backend for BigTableBackend { let path = id.as_storage_path().to_string().into_bytes(); self.mutate(path, [delete_row_mutation()], "delete").await?; - self.change_stream.delete(id); + self.change_stream.commit_delete(id).await; Ok(()) } @@ -1136,6 +1141,9 @@ impl HighVolumeBackend for BigTableBackend { access_time: Timestamp, ) -> Result> { objectstore_log::debug!("Conditional put to Bigtable backend"); + self.change_stream + .begin_write(id, metadata.time_expires) + .await?; let path = id.as_storage_path().to_string().into_bytes(); let (mutations, size) = object_mutations(&path, metadata.clone(), payload.to_vec())?; @@ -1151,7 +1159,9 @@ impl HighVolumeBackend for BigTableBackend { .await?; if write_succeeded { - self.change_stream.write(id, size, metadata.time_expires); + self.change_stream + .commit_write(id, size, metadata.time_expires) + .await; return Ok(None); } @@ -1363,12 +1373,13 @@ impl HighVolumeBackend for BigTableBackend { } }; + self.change_stream.begin_update(id, Some(expire_at)).await?; let applied = self .check_and_mutate(path, predicate, mutations, "set_expiry") .await?; if applied { - self.change_stream.update(id, Some(expire_at)); + self.change_stream.commit_update(id, Some(expire_at)).await; } Ok(if applied { @@ -1399,7 +1410,7 @@ impl HighVolumeBackend for BigTableBackend { .await?; if deleted { - self.change_stream.delete(id); + self.change_stream.commit_delete(id).await; return Ok(None); } @@ -1471,6 +1482,9 @@ impl HighVolumeBackend for BigTableBackend { } TieredWrite::Delete => (vec![delete_row_mutation()], None), }; + if let Some(expires_at) = expires_at { + self.change_stream.begin_write(id, expires_at).await?; + } let written = self .check_and_mutate( @@ -1487,10 +1501,11 @@ impl HighVolumeBackend for BigTableBackend { // We wrote something (the inner `expires_at` is `None` for manual GC) (true, Some(expires_at)) => { self.change_stream - .write(id, row_size(&path, &mutations), expires_at) + .commit_write(id, row_size(&path, &mutations), expires_at) + .await } // We deleted something - (true, None) => self.change_stream.delete(id), + (true, None) => self.change_stream.commit_delete(id).await, } Ok(written) diff --git a/objectstore-service/src/backend/gcs.rs b/objectstore-service/src/backend/gcs.rs index 7b870f92..acb549cb 100644 --- a/objectstore-service/src/backend/gcs.rs +++ b/objectstore-service/src/backend/gcs.rs @@ -741,8 +741,8 @@ impl GcsBackend { .await } - /// Reports an object write to the [`ChangeStream`]. - fn report_object_write( + /// Reports a committed object write to the [`ChangeStream`]. + async fn report_object_write( &self, id: &ObjectId, stored_size: Option, @@ -752,7 +752,8 @@ impl GcsBackend { match stored_size { Some(stored_size) => { self.change_stream - .write(id, stored_size + metadata_size, expires_at) + .commit_write(id, stored_size + metadata_size, expires_at) + .await } None => { objectstore_metrics::count!("change_stream.unreported", reason = "no_stored_size") @@ -884,6 +885,9 @@ impl Backend for GcsBackend { _access_time: Timestamp, ) -> Result { objectstore_log::debug!("Writing to GCS backend"); + self.change_stream + .begin_write(id, metadata.time_expires) + .await?; let gcs_metadata = GcsObject::from_metadata(metadata); // NB: Ensure the order of these fields and that a content-type is attached to them. Both @@ -932,7 +936,8 @@ impl Backend for GcsBackend { stored_size, gcs_metadata.metadata_size(), metadata.time_expires, - ); + ) + .await; Ok(()) } @@ -1096,11 +1101,12 @@ impl Backend for GcsBackend { None }; + self.change_stream.begin_update(id, Some(expire_at)).await?; let outcome = self .update_custom_time(object_url, expire_at, object.generations(), metadata) .await?; if matches!(outcome, SetExpiryResponse::Satisfied(_)) { - self.change_stream.update(id, Some(expire_at)); + self.change_stream.commit_update(id, Some(expire_at)).await; } Ok(outcome) @@ -1140,7 +1146,7 @@ impl Backend for GcsBackend { .await?; if deleted { - self.change_stream.delete(id); + self.change_stream.commit_delete(id).await; } Ok(()) @@ -1154,6 +1160,9 @@ impl Backend for GcsBackend { upload_length: NonZeroU64, ) -> Result> { objectstore_log::debug!("Creating resumable upload session on GCS backend"); + self.change_stream + .begin_write(id, metadata.time_expires) + .await?; let url = self.upload_url(id, "resumable")?; let metadata_json = serde_json::to_vec(&GcsObject::from_metadata(metadata)).context( ErrorKind::Internal, @@ -1269,7 +1278,8 @@ impl Backend for GcsBackend { stored_size, object.metadata_size(), expires_at, - ); + ) + .await; } Ok(progress.into()) } @@ -1304,7 +1314,8 @@ impl Backend for GcsBackend { stored_size, object.metadata_size(), expires_at, - ); + ) + .await; } Ok(progress.into()) }) @@ -1474,6 +1485,9 @@ impl MultipartUploadBackend for GcsBackend { metadata: &Metadata, ) -> Result { objectstore_log::debug!("Initiating multipart upload on GCS backend"); + self.change_stream + .begin_write(id, metadata.time_expires) + .await?; let mut url = self.xml_object_url(id)?; url.set_query(Some("uploads")); @@ -1656,7 +1670,7 @@ impl MultipartUploadBackend for GcsBackend { id: &ObjectId, upload_id: &UploadId, parts: Vec, - _access_time: Timestamp, + access_time: Timestamp, ) -> Result { objectstore_log::debug!("Completing multipart upload on GCS backend"); let mut url = self.xml_object_url(id)?; @@ -1695,6 +1709,27 @@ impl MultipartUploadBackend for GcsBackend { .ok() .map(Into::into); + if error.is_none() { + // The XML API does not return the stored object, so read it back. + let object = self + .get_gcs_metadata(&self.object_url(id)?, access_time) + .await + .ok() + .flatten(); + let stored_size = object + .as_ref() + .and_then(|object| object.size.as_deref()?.parse().ok()); + self.report_object_write( + id, + stored_size, + object.as_ref().map_or(0, GcsObject::metadata_size), + object.and_then(|object| { + object.custom_time.map(Rfc3339Timestamp::into_inner) + }), + ) + .await; + } + Ok(error) } }) @@ -3235,6 +3270,46 @@ mod tests { Ok(()) } + #[cfg(feature = "storage-cogs")] + #[tokio::test] + async fn multipart_completion_reports_to_change_stream() -> Result<()> { + let (backend, producer) = create_test_backend_with_change_stream().await?; + let id = make_id(); + let payload = b"multipart payload".to_vec(); + let metadata = Metadata { + time_expires: Some(Timestamp::now() + Duration::from_secs(3600)), + ..Default::default() + }; + let upload_id = backend.initiate_multipart(&id, &metadata).await?; + let etag = backend + .upload_part( + &id, + &upload_id, + NonZeroU32::new(1).unwrap(), + payload.len() as u64, + None, + stream::single(payload.clone()), + ) + .await?; + let part = CompletedPart { + part_number: NonZeroU32::new(1).unwrap(), + etag, + }; + backend + .complete_multipart(&id, &upload_id, vec![part], Timestamp::now()) + .await?; + + let records = producer.records(); + assert_eq!(records.len(), 1); + assert_eq!(records[0].op_type, OpType::Write); + assert_eq!( + records[0].size, + Some(payload.len() as u64 + GcsObject::from_metadata(&metadata).metadata_size()) + ); + assert!(records[0].expiration_time.is_some()); + Ok(()) + } + #[cfg(feature = "storage-cogs")] #[tokio::test] async fn resumable_completion_reports_to_change_stream() -> Result<()> { diff --git a/objectstore-service/src/backend/in_memory.rs b/objectstore-service/src/backend/in_memory.rs index 4e022b57..1f2a83ff 100644 --- a/objectstore-service/src/backend/in_memory.rs +++ b/objectstore-service/src/backend/in_memory.rs @@ -132,6 +132,45 @@ impl InMemoryBackend { self } + /// Applies `extend` to `id`'s entry, calling [`ChangeStream::begin_update`] first. + /// + /// Like a conditional write, this is rejected if the entry changed in the meantime. + async fn update_expiry( + &self, + id: &ObjectId, + access_time: Timestamp, + extend: impl Fn(&mut StoreEntry) -> Result + Send + Sync, + ) -> Result { + let try_extend = |store: &Store| -> Result<(ExpiryOutcome, Option)> { + match store.get(id) { + Some(entry) if !entry.is_expired(access_time) => { + let mut updated = entry.clone(); + let outcome = extend(&mut updated)?; + Ok((outcome, Some(updated))) + } + _ => Ok((ExpiryOutcome::NotFound, None)), + } + }; + + let (outcome, _) = try_extend(&self.store.lock().unwrap())?; + let ExpiryOutcome::Extended(expire_at) = outcome else { + return Ok(outcome.response()); + }; + self.change_stream.begin_update(id, Some(expire_at)).await?; + + { + let mut store = self.store.lock().unwrap(); + match try_extend(&store)? { + (ExpiryOutcome::Extended(t), Some(updated)) if t == expire_at => { + store.insert(id.clone(), updated); + } + _ => return Ok(SetExpiryResponse::Rejected), + } + } + self.change_stream.commit_update(id, Some(expire_at)).await; + Ok(outcome.response()) + } + /// Returns the stored entry for `id`, for direct inspection in tests. pub fn get(&self, id: &ObjectId) -> Entry { match self.store.lock().unwrap().get(id).cloned() { @@ -176,12 +215,16 @@ impl super::common::Backend for InMemoryBackend { stream: ClientStream, _access_time: Timestamp, ) -> Result { + self.change_stream + .begin_write(id, metadata.time_expires) + .await?; let bytes: BytesMut = stream.try_collect().await?; let entry = StoreEntry::Object(metadata.clone(), bytes.freeze()); let size = entry.stored_size(); self.store.lock().unwrap().insert(id.clone(), entry); self.change_stream - .write(id, size as u64, metadata.time_expires); + .commit_write(id, size as u64, metadata.time_expires) + .await; Ok(()) } @@ -225,23 +268,11 @@ impl super::common::Backend for InMemoryBackend { target: ExpiryUpdate, access_time: Timestamp, ) -> Result { - let outcome = { - let mut store = self.store.lock().unwrap(); - match store.get_mut(id) { - None => ExpiryOutcome::NotFound, - Some(entry) if entry.is_expired(access_time) => ExpiryOutcome::NotFound, - Some(StoreEntry::Object(metadata, _)) => { - extend_object_expiry(metadata, target, access_time)? - } - _ => ExpiryOutcome::Rejected, - } - }; - - if let ExpiryOutcome::Extended(expire_at) = outcome { - self.change_stream.update(id, Some(expire_at)); - } - - Ok(outcome.response()) + self.update_expiry(id, access_time, |entry| match entry { + StoreEntry::Object(metadata, _) => extend_object_expiry(metadata, target, access_time), + StoreEntry::Tombstone(_) => Ok(ExpiryOutcome::Rejected), + }) + .await } async fn delete_object( @@ -250,7 +281,7 @@ impl super::common::Backend for InMemoryBackend { _access_time: Timestamp, ) -> Result { if self.store.lock().unwrap().remove(id).is_some() { - self.change_stream.delete(id); + self.change_stream.commit_delete(id).await; } Ok(()) } @@ -261,6 +292,9 @@ impl super::common::Backend for InMemoryBackend { metadata: &Metadata, _upload_length: NonZeroU64, ) -> Result> { + self.change_stream + .begin_write(id, metadata.time_expires) + .await?; let token = uuid::Uuid::now_v7().to_string(); let upload = ResumableUpload { metadata: metadata.clone(), @@ -324,7 +358,8 @@ impl super::common::Backend for InMemoryBackend { .insert(session.object_id.clone(), entry); self.change_stream - .write(&session.object_id, size as u64, expires_at); + .commit_write(&session.object_id, size as u64, expires_at) + .await; self.resumable_store .lock() .unwrap() @@ -401,20 +436,29 @@ impl HighVolumeBackend for InMemoryBackend { payload: Bytes, access_time: Timestamp, ) -> Result> { - let mut store = self.store.lock().unwrap(); - if let Some(StoreEntry::Tombstone(tombstone)) = store.get(id) - && !tombstone.is_expired(access_time) - { - return Ok(Some(tombstone.clone())); - } + self.change_stream + .begin_write(id, metadata.time_expires) + .await?; + let (size, expires_at) = { + let mut store = self.store.lock().unwrap(); + if let Some(StoreEntry::Tombstone(tombstone)) = store.get(id) + && !tombstone.is_expired(access_time) + { + return Ok(Some(tombstone.clone())); + } - let mut metadata = metadata.clone(); - metadata.size = Some(payload.len()); - let expires_at = metadata.time_expires; - let entry = StoreEntry::Object(metadata, payload); - let size = entry.stored_size(); - store.insert(id.clone(), entry); - self.change_stream.write(id, size as u64, expires_at); + let mut metadata = metadata.clone(); + metadata.size = Some(payload.len()); + let expires_at = metadata.time_expires; + let entry = StoreEntry::Object(metadata, payload); + let size = entry.stored_size(); + store.insert(id.clone(), entry); + (size, expires_at) + }; + + self.change_stream + .commit_write(id, size as u64, expires_at) + .await; Ok(None) } @@ -467,15 +511,19 @@ impl HighVolumeBackend for InMemoryBackend { id: &ObjectId, access_time: Timestamp, ) -> Result> { - let mut store = self.store.lock().unwrap(); - if let Some(StoreEntry::Tombstone(tombstone)) = store.get(id).cloned() - && !tombstone.is_expired(access_time) - { - return Ok(Some(tombstone)); - } + let removed = { + let mut store = self.store.lock().unwrap(); + if let Some(StoreEntry::Tombstone(tombstone)) = store.get(id).cloned() + && !tombstone.is_expired(access_time) + { + return Ok(Some(tombstone)); + } - if store.remove(id).is_some() { - self.change_stream.delete(id); + store.remove(id).is_some() + }; + + if removed { + self.change_stream.commit_delete(id).await; } Ok(None) } @@ -488,26 +536,16 @@ impl HighVolumeBackend for InMemoryBackend { access_time: Timestamp, ) -> Result { let TieredUpdate::SetExpiry(expiry_target) = update; - let outcome = { - let mut store = self.store.lock().unwrap(); - match (store.get_mut(id), current) { - (None, _) => ExpiryOutcome::NotFound, - (Some(entry), _) if entry.is_expired(access_time) => ExpiryOutcome::NotFound, - (Some(StoreEntry::Object(metadata, _)), None) => { - extend_object_expiry(metadata, expiry_target, access_time)? - } - (Some(StoreEntry::Tombstone(t)), Some(target)) if t.target == *target => { - extend_expiry(&mut t.time_expires, expiry_target, None, access_time)? - } - _ => ExpiryOutcome::Rejected, + self.update_expiry(id, access_time, |entry| match (entry, current) { + (StoreEntry::Object(metadata, _), None) => { + extend_object_expiry(metadata, expiry_target, access_time) } - }; - - if let ExpiryOutcome::Extended(expire_at) = outcome { - self.change_stream.update(id, Some(expire_at)); - } - - Ok(outcome.response()) + (StoreEntry::Tombstone(t), Some(target)) if t.target == *target => { + extend_expiry(&mut t.time_expires, expiry_target, None, access_time) + } + _ => Ok(ExpiryOutcome::Rejected), + }) + .await } async fn compare_and_write( @@ -517,37 +555,60 @@ impl HighVolumeBackend for InMemoryBackend { write: TieredWrite, access_time: Timestamp, ) -> Result { - let mut store = self.store.lock().unwrap(); - - let actual = store.get(id); - let matches_current = matches_redirect(actual, current, access_time); - let matches_next = matches_redirect(actual, write.target(), access_time); - - if matches_current { - match write { - TieredWrite::Tombstone(tombstone) => { - let expires_at = tombstone.time_expires; - let entry = StoreEntry::Tombstone(tombstone); - let size = entry.stored_size(); - store.insert(id.clone(), entry); - self.change_stream.write(id, size as u64, expires_at); - } - TieredWrite::Object(metadata, payload) => { - let expires_at = metadata.time_expires; - let entry = StoreEntry::Object(metadata, payload); - let size = entry.stored_size(); - store.insert(id.clone(), entry); - self.change_stream.write(id, size as u64, expires_at); - } - TieredWrite::Delete => { - if store.remove(id).is_some() { - self.change_stream.delete(id); + match &write { + TieredWrite::Tombstone(tombstone) => { + self.change_stream + .begin_write(id, tombstone.time_expires) + .await?; + } + TieredWrite::Object(metadata, _) => { + self.change_stream + .begin_write(id, metadata.time_expires) + .await?; + } + TieredWrite::Delete => {} + } + + let mut written = None; + let mut deleted = false; + let matches = { + let mut store = self.store.lock().unwrap(); + + let actual = store.get(id); + let matches_current = matches_redirect(actual, current, access_time); + let matches_next = matches_redirect(actual, write.target(), access_time); + + if matches_current { + match write { + TieredWrite::Tombstone(tombstone) => { + let expires_at = tombstone.time_expires; + let entry = StoreEntry::Tombstone(tombstone); + written = Some((entry.stored_size(), expires_at)); + store.insert(id.clone(), entry); + } + TieredWrite::Object(metadata, payload) => { + let expires_at = metadata.time_expires; + let entry = StoreEntry::Object(metadata, payload); + written = Some((entry.stored_size(), expires_at)); + store.insert(id.clone(), entry); } + TieredWrite::Delete => deleted = store.remove(id).is_some(), } } + + matches_current || matches_next + }; + + if let Some((size, expires_at)) = written { + self.change_stream + .commit_write(id, size as u64, expires_at) + .await; + } + if deleted { + self.change_stream.commit_delete(id).await; } - Ok(matches_current || matches_next) + Ok(matches) } } @@ -558,6 +619,9 @@ impl MultipartUploadBackend for InMemoryBackend { id: &ObjectId, metadata: &Metadata, ) -> Result { + self.change_stream + .begin_write(id, metadata.time_expires) + .await?; let upload_id = UploadId::new(uuid::Uuid::now_v7().to_string())?; let upload = MultipartUpload { metadata: metadata.clone(), @@ -729,7 +793,9 @@ impl MultipartUploadBackend for InMemoryBackend { let entry = StoreEntry::Object(metadata, payload); let size = entry.stored_size(); self.store.lock().unwrap().insert(id.clone(), entry); - self.change_stream.write(id, size as u64, expires_at); + self.change_stream + .commit_write(id, size as u64, expires_at) + .await; self.multipart_store.lock().unwrap().remove(&key); diff --git a/objectstore-service/src/backend/local_fs.rs b/objectstore-service/src/backend/local_fs.rs index 751c64b1..02809cfd 100644 --- a/objectstore-service/src/backend/local_fs.rs +++ b/objectstore-service/src/backend/local_fs.rs @@ -167,6 +167,9 @@ impl Backend for LocalFsBackend { ) -> Result { let path = self.path(id); objectstore_log::debug!(path=%path.display(), "Writing to local_fs backend"); + self.change_stream + .begin_write(id, metadata.time_expires) + .await?; create_directories(path.parent().unwrap()).await.context( ErrorKind::BackendFailure, "creating local-fs object directory", @@ -192,7 +195,8 @@ impl Backend for LocalFsBackend { draft.publish().await?; self.change_stream - .write(id, stored_size, metadata.time_expires); + .commit_write(id, stored_size, metadata.time_expires) + .await; Ok(()) } @@ -273,6 +277,7 @@ impl Backend for LocalFsBackend { expire_at, )?; metadata.time_expires = Some(expire_at); + self.change_stream.begin_update(id, Some(expire_at)).await?; let mut draft = Draft::create(&path, &metadata).await?; tokio::io::copy(&mut reader, draft.writer()).await.context( @@ -283,7 +288,7 @@ impl Backend for LocalFsBackend { draft.prepare().await?; draft.publish().await?; - self.change_stream.update(id, Some(expire_at)); + self.change_stream.commit_update(id, Some(expire_at)).await; Ok(SetExpiryResponse::Satisfied(expire_at)) } @@ -299,7 +304,7 @@ impl Backend for LocalFsBackend { objectstore_log::debug!("Deleting from local_fs backend"); let path = self.path(id); match tokio::fs::remove_file(path).await { - Ok(()) => self.change_stream.delete(id), + Ok(()) => self.change_stream.commit_delete(id).await, Err(error) if error.kind() == io::ErrorKind::NotFound => { objectstore_log::debug!("Object not found"); } @@ -318,6 +323,9 @@ impl Backend for LocalFsBackend { metadata: &Metadata, _upload_length: NonZeroU64, ) -> Result> { + self.change_stream + .begin_write(id, metadata.time_expires) + .await?; let upload_id = uuid::Uuid::now_v7(); let path = self.upload_path(upload_id); create_directories(path.parent().unwrap()).await.context( @@ -377,7 +385,8 @@ impl Backend for LocalFsBackend { )?; let (stored_size, expires_at) = upload.publish(object_path).await?; self.change_stream - .write(&session.object_id, stored_size, expires_at); + .commit_write(&session.object_id, stored_size, expires_at) + .await; Ok(UploadProgress::Complete) } @@ -426,6 +435,9 @@ impl MultipartUploadBackend for LocalFsBackend { id: &ObjectId, metadata: &Metadata, ) -> Result { + self.change_stream + .begin_write(id, metadata.time_expires) + .await?; let upload_id = UploadId::new(Uuid::now_v7().to_string())?; let dir = self.multipart_dir(id, &upload_id); create_directories(&dir).await.context( @@ -741,7 +753,8 @@ impl MultipartUploadBackend for LocalFsBackend { drop(guard); self.change_stream - .write(id, stored_size, metadata.time_expires); + .commit_write(id, stored_size, metadata.time_expires) + .await; // Clean up multipart state tokio::fs::remove_dir_all(dir).await.context( diff --git a/objectstore-service/src/backend/s3_compatible.rs b/objectstore-service/src/backend/s3_compatible.rs index b0bcac04..bdff4c35 100644 --- a/objectstore-service/src/backend/s3_compatible.rs +++ b/objectstore-service/src/backend/s3_compatible.rs @@ -379,6 +379,9 @@ impl Backend for S3CompatibleBackend { _access_time: Timestamp, ) -> Result { objectstore_log::debug!("Writing to s3_compatible backend"); + self.change_stream + .begin_write(id, metadata.time_expires) + .await?; let headers = metadata_to_gcs_headers(metadata, GCS_CUSTOM_PREFIX) .context(ErrorKind::InvalidMetadata, "encoding S3 object metadata")?; let metadata_size = headers_size(&headers); @@ -397,11 +400,13 @@ impl Backend for S3CompatibleBackend { .drain_body() .await; - self.change_stream.write( - id, - metadata_size + payload_size.load(Ordering::Relaxed), - metadata.time_expires, - ); + self.change_stream + .commit_write( + id, + metadata_size + payload_size.load(Ordering::Relaxed), + metadata.time_expires, + ) + .await; Ok(()) } @@ -477,11 +482,12 @@ impl Backend for S3CompatibleBackend { expire_at, )?; metadata.time_expires = Some(expire_at); + self.change_stream.begin_update(id, Some(expire_at)).await?; let outcome = self .update_metadata(id, &metadata, expire_at, &etag) .await?; if matches!(outcome, SetExpiryResponse::Satisfied(_)) { - self.change_stream.update(id, Some(expire_at)); + self.change_stream.commit_update(id, Some(expire_at)).await; } Ok(outcome) @@ -516,7 +522,7 @@ impl Backend for S3CompatibleBackend { // If the object didn't exist in the first place, this emits a spurious message // due to S3 returning 204 to DELETEs whether the object existed or not. - self.change_stream.delete(id); + self.change_stream.commit_delete(id).await; Ok(()) } diff --git a/objectstore-service/src/change_stream/cost_tracker.rs b/objectstore-service/src/change_stream/cost_tracker.rs index 33d805ce..97341346 100644 --- a/objectstore-service/src/change_stream/cost_tracker.rs +++ b/objectstore-service/src/change_stream/cost_tracker.rs @@ -73,7 +73,7 @@ where P: Producer + Clone + Send + Sync + 'static, P::Error: Into + Send + 'static, { - fn write(&self, id: &ObjectId, size: u64, expires_at: Option) { + async fn commit_write(&self, id: &ObjectId, size: u64, expires_at: Option) { let result = self.tracker.write( &id.as_storage_path().to_string(), id.usecase(), @@ -86,7 +86,7 @@ where self.swallow("write", result); } - fn update(&self, id: &ObjectId, expires_at: Option) { + async fn commit_update(&self, id: &ObjectId, expires_at: Option) { let result = self.tracker.update( &id.as_storage_path().to_string(), id.usecase(), @@ -98,7 +98,7 @@ where self.swallow("update", result); } - fn delete(&self, id: &ObjectId) { + async fn commit_delete(&self, id: &ObjectId) { let result = self.tracker.delete( &id.as_storage_path().to_string(), id.usecase(), @@ -135,12 +135,12 @@ mod tests { (producer, stream) } - #[test] - fn usecase_and_scopes_are_extracted_from_the_id() { + #[tokio::test] + async fn usecase_and_scopes_are_extracted_from_the_id() { let (producer, stream) = stream(1.0); let id = object_id("attachments/org.17/project.42/objects/abc"); - stream.write(&id, 4096, None); + stream.commit_write(&id, 4096, None).await; let record = &producer.records()[0]; assert_eq!(record.shared_resource_id, "bigtable_objectstore"); @@ -151,20 +151,20 @@ mod tests { assert_eq!(record.op_type, OpType::Write); } - #[test] - fn the_storage_path_is_not_emitted() { + #[tokio::test] + async fn the_storage_path_is_not_emitted() { let (producer, stream) = stream(1.0); let id = object_id("attachments/org.17/project.42/objects/abc"); - stream.write(&id, 4096, None); + stream.commit_write(&id, 4096, None).await; let record_id = &producer.records()[0].record_id; assert_ne!(record_id, &id.as_storage_path().to_string()); assert!(!record_id.contains("attachments")); } - #[test] - fn missing_or_unparseable_scopes_are_reported_as_absent() { + #[tokio::test] + async fn missing_or_unparseable_scopes_are_reported_as_absent() { let (producer, stream) = stream(1.0); for path in [ @@ -172,7 +172,7 @@ mod tests { "attachments/organization.17/objects/abc", "attachments/org.not-a-number/project.42/objects/abc", ] { - stream.write(&object_id(path), 1, None); + stream.commit_write(&object_id(path), 1, None).await; } let records = producer.records(); @@ -200,14 +200,14 @@ mod tests { } } - #[test] - fn every_operation_on_an_object_reports_the_same_record() { + #[tokio::test] + async fn every_operation_on_an_object_reports_the_same_record() { let (producer, stream) = stream(1.0); let id = object_id("attachments/org.1/project.2/objects/abc"); - stream.write(&id, 10, None); - stream.update(&id, Some(Timestamp::now())); - stream.delete(&id); + stream.commit_write(&id, 10, None).await; + stream.commit_update(&id, Some(Timestamp::now())).await; + stream.commit_delete(&id).await; let records = producer.records(); assert_eq!(records.len(), 3); @@ -215,33 +215,37 @@ mod tests { assert_eq!(records[1].record_id, records[2].record_id); } - #[test] - fn distinct_revisions_are_distinct_records() { + #[tokio::test] + async fn distinct_revisions_are_distinct_records() { let (producer, stream) = stream(1.0); - stream.write( - &object_id("attachments/org.1/project.2/objects/abc/0199aaaa"), - 1, - None, - ); - stream.write( - &object_id("attachments/org.1/project.2/objects/abc/0199bbbb"), - 1, - None, - ); + stream + .commit_write( + &object_id("attachments/org.1/project.2/objects/abc/0199aaaa"), + 1, + None, + ) + .await; + stream + .commit_write( + &object_id("attachments/org.1/project.2/objects/abc/0199bbbb"), + 1, + None, + ) + .await; let records = producer.records(); assert_ne!(records[0].record_id, records[1].record_id); } - #[test] - fn update_omits_size_and_delete_omits_everything_optional() { + #[tokio::test] + async fn update_omits_size_and_delete_omits_everything_optional() { let (producer, stream) = stream(1.0); let id = object_id("attachments/org.1/project.2/objects/abc"); let expires = Timestamp::from_unix_micros(1_800_000_000_123_456).unwrap(); - stream.update(&id, Some(expires)); - stream.delete(&id); + stream.commit_update(&id, Some(expires)).await; + stream.commit_delete(&id).await; let records = producer.records(); assert_eq!(records[0].op_type, OpType::Update); @@ -252,15 +256,15 @@ mod tests { assert_eq!(records[1].expiration_time, None); } - #[test] - fn a_listener_sampled_at_zero_reports_nothing() { + #[tokio::test] + async fn a_listener_sampled_at_zero_reports_nothing() { let (producer, stream) = stream(0.0); for i in 0..100 { let id = object_id(&format!("attachments/org.1/project.2/objects/{i}")); - stream.write(&id, 1, None); - stream.update(&id, None); - stream.delete(&id); + stream.commit_write(&id, 1, None).await; + stream.commit_update(&id, None).await; + stream.commit_delete(&id).await; } assert!(producer.records().is_empty()); diff --git a/objectstore-service/src/change_stream/factory.rs b/objectstore-service/src/change_stream/factory.rs index d259c142..648d962b 100644 --- a/objectstore-service/src/change_stream/factory.rs +++ b/objectstore-service/src/change_stream/factory.rs @@ -149,14 +149,18 @@ mod tests { assert!(!reports(&factory.build(Some(&config())))); } - #[test] - fn a_configured_backend_reports_through_the_transport() { + #[tokio::test] + async fn a_configured_backend_reports_through_the_transport() { let (factory, producer) = dummy_factory(); let stream = factory.build(Some(&config())); assert!(reports(&stream)); - stream.delete(&crate::id::ObjectId::from_storage_path("attachments/objects/abc").unwrap()); + stream + .commit_delete( + &crate::id::ObjectId::from_storage_path("attachments/objects/abc").unwrap(), + ) + .await; let records = producer.records(); assert_eq!(records.len(), 1); diff --git a/objectstore-service/src/change_stream/mod.rs b/objectstore-service/src/change_stream/mod.rs index 0c9473f5..37f11790 100644 --- a/objectstore-service/src/change_stream/mod.rs +++ b/objectstore-service/src/change_stream/mod.rs @@ -77,17 +77,55 @@ fn default_sample_rate() -> f64 { /// Publishes the changes a single backend makes to the objects it stores. /// +/// Backends call `begin_*` before a storage operation that creates data or extends its +/// lifetime, and `commit_*` after a storage operation succeeds. A failed `begin_*` fails +/// the request before storage is touched. A `begin_*` is not always followed by a +/// `commit_*` (e.g. when the storage operation fails), so consumers that must never +/// underestimate what is stored should treat `begin_*` as an upper bound. +/// +/// All methods but [`join`](Self::join) default to doing nothing. Backends await them +/// on the request path, so they add to request latency. +/// +/// Consider spawning and independent task or using a queue in your implementation of the +/// methods if you only need best-effort delivery. +/// /// See [module docs](self). #[async_trait::async_trait] pub trait ChangeStream: fmt::Debug + Send + Sync + 'static { + /// Announces that `id` is about to be written. Used for new writes and overwrites. + async fn begin_write( + &self, + id: &ObjectId, + expires_at: Option, + ) -> Result<(), ChangeStreamError> { + let _ = (id, expires_at); + Ok(()) + } + /// Reports that `id` now occupies `size` bytes. Used for new writes and overwrites. - fn write(&self, id: &ObjectId, size: u64, expires_at: Option); + async fn commit_write(&self, id: &ObjectId, size: u64, expires_at: Option) { + let _ = (id, size, expires_at); + } + + /// Announces that `id`'s expiration is about to move out to `expires_at`. + async fn begin_update( + &self, + id: &ObjectId, + expires_at: Option, + ) -> Result<(), ChangeStreamError> { + let _ = (id, expires_at); + Ok(()) + } /// Reports that `id`'s expiration moved, with its stored size unchanged. - fn update(&self, id: &ObjectId, expires_at: Option); + async fn commit_update(&self, id: &ObjectId, expires_at: Option) { + let _ = (id, expires_at); + } /// Reports that `id` was deleted explicitly. Does not account for automatic GC. - fn delete(&self, id: &ObjectId); + async fn commit_delete(&self, id: &ObjectId) { + let _ = id; + } /// Blocks until reported records have been delivered, or `timeout` elapses. /// @@ -99,6 +137,30 @@ pub trait ChangeStream: fmt::Debug + Send + Sync + 'static { async fn join(&self, timeout: Duration); } +/// Why a [`ChangeStream`] could not begin a change. +/// +/// The variant determines the [`ErrorKind`](crate::error::ErrorKind) of the failed +/// request. The source carries the details and can be a plain message: +/// +/// ``` +/// use objectstore_service::change_stream::ChangeStreamError; +/// +/// let error = ChangeStreamError::Failure("constraint violated".into()); +/// assert_eq!(error.to_string(), "change stream failed"); +/// ``` +#[derive(Debug, thiserror::Error)] +pub enum ChangeStreamError { + /// The sink is temporarily unavailable. + #[error("change stream unavailable")] + Unavailable(#[source] Box), + /// The sink did not respond in time. + #[error("change stream timed out")] + Timeout(#[source] Box), + /// Any other failure. + #[error("change stream failed")] + Failure(#[source] Box), +} + /// Drains `change_stream`, bounded by [`FLUSH_TIMEOUT`]. /// /// Backends call this from [`Backend::join`](crate::backend::common::Backend::join) so @@ -113,12 +175,6 @@ pub struct NoopStream; #[async_trait::async_trait] impl ChangeStream for NoopStream { - fn write(&self, _id: &ObjectId, _size: u64, _expires_at: Option) {} - - fn update(&self, _id: &ObjectId, _expires_at: Option) {} - - fn delete(&self, _id: &ObjectId) {} - async fn join(&self, _timeout: Duration) {} } diff --git a/objectstore-service/src/error.rs b/objectstore-service/src/error.rs index b38ed3f6..8faf2904 100644 --- a/objectstore-service/src/error.rs +++ b/objectstore-service/src/error.rs @@ -10,6 +10,9 @@ use std::error::Error as StdError; use std::fmt; use objectstore_log::Level; + +use crate::change_stream::ChangeStreamError; + /// A panic captured from a service task. #[derive(Debug)] pub struct Panic { @@ -311,6 +314,17 @@ where } } +impl From for Error { + fn from(source: ChangeStreamError) -> Self { + let kind = match source { + ChangeStreamError::Unavailable(_) => ErrorKind::BackendUnavailable, + ChangeStreamError::Timeout(_) => ErrorKind::BackendTimeout, + ChangeStreamError::Failure(_) => ErrorKind::BackendFailure, + }; + Self::with_context(kind, "writing to the change stream", source) + } +} + impl From for Error { fn from(source: std::io::Error) -> Self { Self::with_source(ErrorKind::BackendFailure, source)