Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 9 additions & 2 deletions objectstore-service/docs/architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -182,15 +182,22 @@ 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.
- `update(id, expires_at)`: `id`'s expiration moved while its stored size is
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
Expand Down
29 changes: 22 additions & 7 deletions objectstore-service/src/backend/bigtable.rs
Original file line number Diff line number Diff line change
Expand Up @@ -999,6 +999,9 @@ impl Backend for BigTableBackend {
_access_time: Timestamp,
) -> Result<PutResponse> {
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);
Expand All @@ -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(())
}
Expand Down Expand Up @@ -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(())
}
Expand Down Expand Up @@ -1136,6 +1141,9 @@ impl HighVolumeBackend for BigTableBackend {
access_time: Timestamp,
) -> Result<Option<Tombstone>> {
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())?;
Expand All @@ -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);
}

Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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);
}

Expand Down Expand Up @@ -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(
Expand All @@ -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)
Expand Down
93 changes: 84 additions & 9 deletions objectstore-service/src/backend/gcs.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<u64>,
Expand All @@ -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")
Expand Down Expand Up @@ -884,6 +885,9 @@ impl Backend for GcsBackend {
_access_time: Timestamp,
) -> Result<PutResponse> {
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
Expand Down Expand Up @@ -932,7 +936,8 @@ impl Backend for GcsBackend {
stored_size,
gcs_metadata.metadata_size(),
metadata.time_expires,
);
)
.await;

Ok(())
}
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -1140,7 +1146,7 @@ impl Backend for GcsBackend {
.await?;

if deleted {
self.change_stream.delete(id);
self.change_stream.commit_delete(id).await;
}

Ok(())
Expand All @@ -1154,6 +1160,9 @@ impl Backend for GcsBackend {
upload_length: NonZeroU64,
) -> Result<Option<BackendToken>> {
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,
Expand Down Expand Up @@ -1269,7 +1278,8 @@ impl Backend for GcsBackend {
stored_size,
object.metadata_size(),
expires_at,
);
)
.await;
}
Ok(progress.into())
}
Expand Down Expand Up @@ -1304,7 +1314,8 @@ impl Backend for GcsBackend {
stored_size,
object.metadata_size(),
expires_at,
);
)
.await;
}
Ok(progress.into())
})
Expand Down Expand Up @@ -1474,6 +1485,9 @@ impl MultipartUploadBackend for GcsBackend {
metadata: &Metadata,
) -> Result<InitiateMultipartResponse> {
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"));

Expand Down Expand Up @@ -1656,7 +1670,7 @@ impl MultipartUploadBackend for GcsBackend {
id: &ObjectId,
upload_id: &UploadId,
parts: Vec<CompletedPart>,
_access_time: Timestamp,
access_time: Timestamp,
) -> Result<CompleteMultipartResponse> {
objectstore_log::debug!("Completing multipart upload on GCS backend");
let mut url = self.xml_object_url(id)?;
Expand Down Expand Up @@ -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;
}

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Multipart commit skips expired objects

Medium Severity

After a successful multipart complete, size and expiry are taken from get_gcs_metadata, which returns none for expired objects and whose errors are discarded. report_object_write then skips commit_write, so a stored object is omitted from the change stream. Long uploads that outlive their TTL are the typical trigger; put_object does not have this gap because it reads size from the upload response.

Fix in Cursor Fix in Web

Reviewed by Cursor Bugbot for commit 7fd1b64. Configure here.


Ok(error)
}
})
Expand Down Expand Up @@ -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<()> {
Expand Down
Loading
Loading