diff --git a/objectstore-service/src/backend/bigtable.rs b/objectstore-service/src/backend/bigtable.rs index 04df7e5c..cf101725 100644 --- a/objectstore-service/src/backend/bigtable.rs +++ b/objectstore-service/src/backend/bigtable.rs @@ -1086,12 +1086,10 @@ impl HighVolumeBackend for BigTableBackend { timestamp_micros: time_expires.as_micros() as i64, value: vec![1], }))]; - self.mutate( - revision.as_upload_path().to_string().into_bytes(), - mutations, - "create_upload_marker", - ) - .await?; + let path = revision.as_upload_path().to_string().into_bytes(); + let size = row_size(&path, &mutations); + self.mutate(path, mutations, "create_upload_marker").await?; + self.change_stream.write(revision, size, Some(time_expires)); Ok(()) } @@ -1118,13 +1116,21 @@ impl HighVolumeBackend for BigTableBackend { revision: &ObjectId, access_time: Timestamp, ) -> Result { - self.check_and_mutate( - revision.as_upload_path().to_string().into_bytes(), - MutatePredicate::Include(live_row_filter(column_filter(COLUMN_UPLOAD), access_time)), - vec![delete_row_mutation()], - "delete_upload_marker", - ) - .await + let deleted = self + .check_and_mutate( + revision.as_upload_path().to_string().into_bytes(), + MutatePredicate::Include(live_row_filter( + column_filter(COLUMN_UPLOAD), + access_time, + )), + vec![delete_row_mutation()], + "delete_upload_marker", + ) + .await?; + if deleted { + self.change_stream.delete(revision); + } + Ok(deleted) } #[tracing::instrument(level = "debug", fields(?id), skip_all)] diff --git a/objectstore-service/src/backend/common.rs b/objectstore-service/src/backend/common.rs index 560c095a..de87d6e6 100644 --- a/objectstore-service/src/backend/common.rs +++ b/objectstore-service/src/backend/common.rs @@ -405,6 +405,8 @@ pub trait HighVolumeBackend: Backend { /// /// `revision` is the upload's unique LT revision. The fixed deadline is independent /// of object expiration and must not be refreshed by subsequent requests. + /// Reports a change-stream write keyed by `revision`, with the marker's stored size + /// and deadline. async fn create_upload_marker( &self, revision: &ObjectId, @@ -418,6 +420,7 @@ pub trait HighVolumeBackend: Backend { /// /// Returns `true` only when this call deletes a live marker. /// Missing and expired markers return `false`, including on a repeated deletion. + /// Reports a change-stream delete keyed by `revision` only when a live marker is deleted. async fn delete_upload_marker( &self, revision: &ObjectId, diff --git a/objectstore-service/src/backend/in_memory.rs b/objectstore-service/src/backend/in_memory.rs index 4e022b57..9d86a8da 100644 --- a/objectstore-service/src/backend/in_memory.rs +++ b/objectstore-service/src/backend/in_memory.rs @@ -369,6 +369,7 @@ impl HighVolumeBackend for InMemoryBackend { .lock() .unwrap() .insert(revision.clone(), time_expires); + self.change_stream.write(revision, 1, Some(time_expires)); Ok(()) } @@ -386,12 +387,16 @@ impl HighVolumeBackend for InMemoryBackend { revision: &ObjectId, access_time: Timestamp, ) -> Result { - Ok(self + let deleted = self .upload_markers .lock() .unwrap() .remove(revision) - .is_some_and(|expiry| expiry >= access_time)) + .is_some_and(|expiry| expiry >= access_time); + if deleted { + self.change_stream.delete(revision); + } + Ok(deleted) } async fn put_non_tombstone(