Skip to content
Closed
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
32 changes: 19 additions & 13 deletions objectstore-service/src/backend/bigtable.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(())
}

Expand All @@ -1118,13 +1116,21 @@ impl HighVolumeBackend for BigTableBackend {
revision: &ObjectId,
access_time: Timestamp,
) -> Result<bool> {
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)]
Expand Down
3 changes: 3 additions & 0 deletions objectstore-service/src/backend/common.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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,
Expand Down
9 changes: 7 additions & 2 deletions objectstore-service/src/backend/in_memory.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(())
}

Expand All @@ -386,12 +387,16 @@ impl HighVolumeBackend for InMemoryBackend {
revision: &ObjectId,
access_time: Timestamp,
) -> Result<bool> {
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(
Expand Down
Loading