From 35ae9763c3474892cde85844ddb33a5d518069c8 Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Wed, 7 Oct 2026 13:35:10 +0200 Subject: [PATCH 01/11] feat(cogs): Account for resumable upload sessions Track upload sessions and HighVolume markers separately from published objects. Report advertised upload sizes until completion, cancellation, or accounting expiration, preserving existing object identities and sampling. Refs FS-563 --- objectstore-service/docs/architecture.md | 39 ++++-- objectstore-service/src/backend/bigtable.rs | 55 +++++--- objectstore-service/src/backend/gcs.rs | 77 +++++++++-- objectstore-service/src/backend/in_memory.rs | 98 +++++++++++--- objectstore-service/src/backend/local_fs.rs | 99 +++++++++++--- .../src/backend/s3_compatible.rs | 6 +- objectstore-service/src/backend/tiered.rs | 84 +++++++++++- .../src/change_stream/cost_tracker.rs | 121 ++++++++++++++---- .../src/change_stream/factory.rs | 4 +- objectstore-service/src/change_stream/mod.rs | 86 +++++++++++-- objectstore-service/src/id.rs | 13 ++ 11 files changed, 553 insertions(+), 129 deletions(-) diff --git a/objectstore-service/docs/architecture.md b/objectstore-service/docs/architecture.md index f195610f..2d7e5d9d 100644 --- a/objectstore-service/docs/architecture.md +++ b/objectstore-service/docs/architecture.md @@ -5,7 +5,7 @@ the `objectstore-server`. # Cargo features -- `storage-cogs`: support for publishing per-object change streams to Kafka for +- `storage-cogs`: support for publishing object and upload-session change streams to Kafka for storage cost attribution. Off by default; it adds a build step to compile `librdkafka` and requires toolchain components we don't otherwise need. Local and sandbox builds don't have a Kafka topic/consumer anyway. @@ -148,32 +148,38 @@ rate-limiting failures at a higher layer) are not counted. This is gated behind the `storage-cogs` Cargo feature. Each backend reports every write/overwrite, applied expiry extension, and delete -it performs on stored objects to a [`ChangeStream`](change_stream::ChangeStream) +it performs on stored objects and upload sessions to a [`ChangeStream`](change_stream::ChangeStream) (see [the change stream section](#change-streams)). To turn this change stream into COGS data, a stream consumer has to merge each change event into an -external table to update an inventory of objects. The inventory table can be +external table to update a storage inventory. The inventory table can be queried to break down each backend's storage utilization by `app_feature`. To enable storage COGS, enable the `storage-cogs` Cargo feature and provide a [`CostTrackerConfig`](change_stream::CostTrackerConfig) for service-wide sink connection details and a [`CostTrackerStreamConfig`](change_stream::CostTrackerStreamConfig) for per-backend information. -Each row in the inventory table has an anonymized hash of an `ObjectId` as well +Each row in the inventory table has an anonymized hash of its change target identity as well as the row's size, expiry, Sentry org/project, `app_feature`, and relevant backend. When using [`TieredStorage`](backend::tiered::TieredStorage)'s long-term backend the inventory table will contain _two rows_ for an object: a row for the actual object and its size in long-term backend, and a separate row for the tombstone and the tombstone's size in the high-volume backend. +While a resumable upload is in progress, separate session rows account for its +advertised size in the upload backend and its marker's stored size in high-volume +storage. See the [change stream module](change_stream) for their lifetimes. Because the change stream does not observe automatic garbage collection, expired -objects must be filtered out when querying the inventory table. +records must be filtered out when querying the inventory table. Under the hood, [`CostTrackerStream`](change_stream::CostTrackerStream) uses [`InventoryTracker`](objectstore_inventory_tracker::InventoryTracker) to publish change events; it is generic over the transport rather than tied to Kafka. Each backend has its own sampling rate to lessen the load put on the stream -processor. Sampling decisions are made -based on [`ObjectId`](id::ObjectId). Each change event includes the sampling rate that was in +processor. For an unchanged backend sample rate, sampling decisions are consistent for +each target identity: sessions have +separate identities from published objects and other sessions for the same object. +See [`CostTrackerStream`](change_stream::CostTrackerStream) for identity details. +Each change event includes the sampling rate that was in effect at the time so that consumers can smooth over the effects of changing the sampling rate. When aggregating, divide each row's value by its `sample_rate`. @@ -181,23 +187,30 @@ See also: [`objectstore_inventory_tracker`] documentation. # Change Streams -Every backend publishes the changes it makes to the objects it stores as a +Every backend publishes the changes it makes to the objects and upload sessions it stores as a [`ChangeStream`](change_stream::ChangeStream). It is a fire-and-forget, per-backend feed of three operations: -- `write(id, size, expires_at)`: `id` now occupies `size` bytes. Used for both +- `write(target, size, expires_at)`: `target` now occupies `size` bytes. Used for both new objects and overwrites. -- `update(id, expires_at)`: `id`'s expiration moved while its stored size is +- `update(target, expires_at)`: `target`'s expiration moved while its stored size is unchanged. In practice this is a TTI bump. -- `delete(id)`: `id` was deleted explicitly. +- `delete(target)`: `target` was deleted explicitly. + +[`ChangeTarget`](change_stream::ChangeTarget) identifies either an object or an upload +session. An upload session carries its object's identity for attribution and a stable +session ID; [`Session`](resumable::Session) converts directly into a session target. 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 storage as well as for the tombstone record in high-volume storage. -`size` is a count of bytes that the backend actually stores for an object. This -includes object payloads, metadata, and sometimes backend-specific overhead. +For objects and markers, `size` is a count of bytes that the backend actually stores. +This includes object payloads, metadata, and sometimes backend-specific overhead. + +See the [change stream module](change_stream) for upload-session accounting, expiration, +and lifecycle reporting. Decorators such as [`CountingBackend`](backend::counting::CountingBackend) and [`TieredStorage`](backend::tiered::TieredStorage) don't publish change streams diff --git a/objectstore-service/src/backend/bigtable.rs b/objectstore-service/src/backend/bigtable.rs index 04df7e5c..cc943833 100644 --- a/objectstore-service/src/backend/bigtable.rs +++ b/objectstore-service/src/backend/bigtable.rs @@ -54,7 +54,7 @@ use crate::backend::common::{ Tombstone, }; use crate::change_stream::{ - ChangeStream, ChangeStreamFactory, CostTrackerStreamConfig, flush_change_stream, + ChangeStream, ChangeStreamFactory, ChangeTarget, CostTrackerStreamConfig, flush_change_stream, }; use crate::error::{Error, ErrorKind, Result, ResultExt as _}; use crate::gcp_auth::PrefetchingTokenProvider; @@ -1009,7 +1009,8 @@ 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 + .write(id.into(), size, metadata.time_expires); Ok(()) } @@ -1063,7 +1064,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.delete(id.into()); Ok(()) } @@ -1086,12 +1087,14 @@ 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( + ChangeTarget::upload_marker(revision), + size, + Some(time_expires), + ); Ok(()) } @@ -1118,13 +1121,22 @@ 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(ChangeTarget::upload_marker(revision)); + } + Ok(deleted) } #[tracing::instrument(level = "debug", fields(?id), skip_all)] @@ -1151,7 +1163,8 @@ impl HighVolumeBackend for BigTableBackend { .await?; if write_succeeded { - self.change_stream.write(id, size, metadata.time_expires); + self.change_stream + .write(id.into(), size, metadata.time_expires); return Ok(None); } @@ -1368,7 +1381,7 @@ impl HighVolumeBackend for BigTableBackend { .await?; if applied { - self.change_stream.update(id, Some(expire_at)); + self.change_stream.update(id.into(), Some(expire_at)); } Ok(if applied { @@ -1399,7 +1412,7 @@ impl HighVolumeBackend for BigTableBackend { .await?; if deleted { - self.change_stream.delete(id); + self.change_stream.delete(id.into()); return Ok(None); } @@ -1487,10 +1500,10 @@ 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) + .write(id.into(), row_size(&path, &mutations), expires_at) } // We deleted something - (true, None) => self.change_stream.delete(id), + (true, None) => self.change_stream.delete(id.into()), } Ok(written) diff --git a/objectstore-service/src/backend/gcs.rs b/objectstore-service/src/backend/gcs.rs index 7b870f92..5bc67633 100644 --- a/objectstore-service/src/backend/gcs.rs +++ b/objectstore-service/src/backend/gcs.rs @@ -24,7 +24,8 @@ use crate::backend::common::{ }; use crate::backend::extensions::{ReqwestResultExt, ResponseExt, SendTraced}; use crate::change_stream::{ - ChangeStream, ChangeStreamFactory, CostTrackerStreamConfig, flush_change_stream, + ChangeStream, ChangeStreamFactory, ChangeTarget, CostTrackerStreamConfig, UPLOAD_SESSION_TTL, + flush_change_stream, }; use crate::error::{Error, ErrorKind, Result, ResultExt as _}; use crate::gcp_auth::PrefetchingTokenProvider; @@ -752,7 +753,7 @@ impl GcsBackend { match stored_size { Some(stored_size) => { self.change_stream - .write(id, stored_size + metadata_size, expires_at) + .write(id.into(), stored_size + metadata_size, expires_at) } None => { objectstore_metrics::count!("change_stream.unreported", reason = "no_stored_size") @@ -1100,7 +1101,7 @@ impl Backend for GcsBackend { .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.update(id.into(), Some(expire_at)); } Ok(outcome) @@ -1140,7 +1141,7 @@ impl Backend for GcsBackend { .await?; if deleted { - self.change_stream.delete(id); + self.change_stream.delete(id.into()); } Ok(()) @@ -1211,7 +1212,16 @@ impl Backend for GcsBackend { "invalid Location URL in GCS resumable upload creation response", ) })?; - Ok(Some(session_uri.into())) + let token = String::from(session_uri); + self.change_stream.write( + ChangeTarget::UploadSession { + object_id: id, + session_id: &token, + }, + upload_length.get(), + Some(Timestamp::now() + UPLOAD_SESSION_TTL), + ); + Ok(Some(token)) } #[tracing::instrument(level = "debug", fields(?session, offset, content_length), skip_all)] @@ -1270,6 +1280,7 @@ impl Backend for GcsBackend { object.metadata_size(), expires_at, ); + self.change_stream.delete(session.into()); } Ok(progress.into()) } @@ -1305,6 +1316,7 @@ impl Backend for GcsBackend { object.metadata_size(), expires_at, ); + self.change_stream.delete(session.into()); } Ok(progress.into()) }) @@ -1350,7 +1362,9 @@ impl Backend for GcsBackend { } } }) - .await + .await?; + self.change_stream.delete(session.into()); + Ok(()) } async fn join(&self) { @@ -3245,10 +3259,23 @@ mod tests { time_expires: Some(Timestamp::now() + Duration::from_secs(3600)), ..Default::default() }; + let before = Timestamp::now(); let token = backend .create_upload_session(&id, &metadata, nonzero(payload.len() as u64)) .await?; + let created = producer.records(); + assert_eq!(created.len(), 1); + assert_eq!(created[0].size, Some(payload.len() as u64)); + let expiration = created[0].expiration_time.unwrap() as u64; + assert!( + ((before + UPLOAD_SESSION_TTL).as_micros() + ..=(Timestamp::now() + UPLOAD_SESSION_TTL).as_micros()) + .contains(&expiration) + ); + backend.upload_offset(&token).await?; + assert_eq!(producer.records().len(), 1); + assert_eq!( backend .put_chunk( @@ -3262,13 +3289,32 @@ mod tests { ); let records = producer.records(); - assert_eq!(records.len(), 1); - assert_eq!(records[0].op_type, OpType::Write); + assert_eq!(records.len(), 3); assert_eq!( - records[0].size, + records.iter().map(|r| r.op_type).collect::>(), + [OpType::Write, OpType::Write, OpType::Delete] + ); + assert_eq!(records[0].record_id, records[2].record_id); + assert_ne!(records[0].record_id, records[1].record_id); + assert_eq!( + records[1].size, Some(payload.len() as u64 + GcsObject::from_metadata(&metadata).metadata_size()) ); - assert!(records[0].expiration_time.is_some()); + assert_eq!( + records[1].expiration_time, + metadata.time_expires.map(|t| t.as_micros() as i64) + ); + producer.clear(); + let canceled = backend + .create_upload_session(&id, &metadata, nonzero(10)) + .await?; + backend.cancel_upload(&canceled).await?; + backend.cancel_upload(&canceled).await?; + let records = producer.records(); + assert_eq!(records.len(), 3); + assert_eq!(records[1].op_type, OpType::Delete); + assert_eq!(records[2].op_type, OpType::Delete); + assert!(records.iter().all(|r| r.record_id == records[0].record_id)); Ok(()) } @@ -3306,10 +3352,15 @@ mod tests { ); let records = producer.records(); - assert_eq!(records.len(), 1); - assert_eq!(records[0].op_type, OpType::Write); + assert_eq!(records.len(), 3); assert_eq!( - records[0].size, + records.iter().map(|r| r.op_type).collect::>(), + [OpType::Write, OpType::Write, OpType::Delete] + ); + assert_eq!(records[0].record_id, records[2].record_id); + assert_ne!(records[0].record_id, records[1].record_id); + assert_eq!( + records[1].size, Some(payload.len() as u64 + GcsObject::from_metadata(&metadata).metadata_size()) ); Ok(()) diff --git a/objectstore-service/src/backend/in_memory.rs b/objectstore-service/src/backend/in_memory.rs index 4e022b57..6650f349 100644 --- a/objectstore-service/src/backend/in_memory.rs +++ b/objectstore-service/src/backend/in_memory.rs @@ -22,7 +22,9 @@ use crate::backend::common::{ PutResponse, SetExpiryResponse, TieredGet, TieredMetadata, TieredUpdate, TieredWrite, Tombstone, }; -use crate::change_stream::{ChangeStream, NoopStream, flush_change_stream}; +use crate::change_stream::{ + ChangeStream, ChangeTarget, NoopStream, UPLOAD_SESSION_TTL, flush_change_stream, +}; use crate::error::{Error, ErrorKind, Result}; use crate::id::ObjectId; use crate::multipart::{ @@ -181,7 +183,7 @@ impl super::common::Backend for InMemoryBackend { let size = entry.stored_size(); self.store.lock().unwrap().insert(id.clone(), entry); self.change_stream - .write(id, size as u64, metadata.time_expires); + .write(id.into(), size as u64, metadata.time_expires); Ok(()) } @@ -238,7 +240,7 @@ impl super::common::Backend for InMemoryBackend { }; if let ExpiryOutcome::Extended(expire_at) = outcome { - self.change_stream.update(id, Some(expire_at)); + self.change_stream.update(id.into(), Some(expire_at)); } Ok(outcome.response()) @@ -250,7 +252,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.delete(id.into()); } Ok(()) } @@ -259,7 +261,7 @@ impl super::common::Backend for InMemoryBackend { &self, id: &ObjectId, metadata: &Metadata, - _upload_length: NonZeroU64, + upload_length: NonZeroU64, ) -> Result> { let token = uuid::Uuid::now_v7().to_string(); let upload = ResumableUpload { @@ -270,6 +272,14 @@ impl super::common::Backend for InMemoryBackend { (id.clone(), token.clone()), Arc::new(tokio::sync::Mutex::new(Some(upload))), ); + self.change_stream.write( + ChangeTarget::UploadSession { + object_id: id, + session_id: &token, + }, + upload_length.get(), + Some(Timestamp::now() + UPLOAD_SESSION_TTL), + ); Ok(Some(token)) } @@ -324,11 +334,12 @@ impl super::common::Backend for InMemoryBackend { .insert(session.object_id.clone(), entry); self.change_stream - .write(&session.object_id, size as u64, expires_at); + .write((&session.object_id).into(), size as u64, expires_at); self.resumable_store .lock() .unwrap() .remove(&(session.object_id.clone(), session.backend_token.clone())); + self.change_stream.delete(session.into()); Ok(UploadProgress::Complete) } @@ -350,6 +361,7 @@ impl super::common::Backend for InMemoryBackend { .lock() .unwrap() .remove(&(session.object_id.clone(), session.backend_token.clone())); + self.change_stream.delete(session.into()); Ok(()) } @@ -369,6 +381,8 @@ impl HighVolumeBackend for InMemoryBackend { .lock() .unwrap() .insert(revision.clone(), time_expires); + self.change_stream + .write(ChangeTarget::upload_marker(revision), 1, Some(time_expires)); Ok(()) } @@ -386,12 +400,17 @@ 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(ChangeTarget::upload_marker(revision)); + } + Ok(deleted) } async fn put_non_tombstone( @@ -414,7 +433,7 @@ impl HighVolumeBackend for InMemoryBackend { 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); + self.change_stream.write(id.into(), size as u64, expires_at); Ok(None) } @@ -475,7 +494,7 @@ impl HighVolumeBackend for InMemoryBackend { } if store.remove(id).is_some() { - self.change_stream.delete(id); + self.change_stream.delete(id.into()); } Ok(None) } @@ -504,7 +523,7 @@ impl HighVolumeBackend for InMemoryBackend { }; if let ExpiryOutcome::Extended(expire_at) = outcome { - self.change_stream.update(id, Some(expire_at)); + self.change_stream.update(id.into(), Some(expire_at)); } Ok(outcome.response()) @@ -530,18 +549,18 @@ impl HighVolumeBackend for InMemoryBackend { 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); + self.change_stream.write(id.into(), 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); + self.change_stream.write(id.into(), size as u64, expires_at); } TieredWrite::Delete => { if store.remove(id).is_some() { - self.change_stream.delete(id); + self.change_stream.delete(id.into()); } } } @@ -729,7 +748,7 @@ 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.write(id.into(), size as u64, expires_at); self.multipart_store.lock().unwrap().remove(&key); @@ -1118,26 +1137,63 @@ mod tests { #[cfg(feature = "storage-cogs")] #[tokio::test] - async fn resumable_emits_only_publication() { + async fn resumable_reports_session_lifecycle() { let (backend, producer) = backend_with_change_stream(); let id = make_id(); + let before = Timestamp::now(); let token = create_session(&backend, &id, 2).await; + let created = producer.records(); + assert_eq!(created.len(), 1); + assert_eq!(created[0].size, Some(2)); + let expiration = created[0].expiration_time.unwrap() as u64; + assert!( + ((before + UPLOAD_SESSION_TTL).as_micros() + ..=(Timestamp::now() + UPLOAD_SESSION_TTL).as_micros()) + .contains(&expiration) + ); - // Partial session state is not reported as a stored object. + // Partial chunks and offset queries do not change the optimistic size. backend .put_chunk(&token, 0, 1, stream::single("a")) .await .unwrap(); - assert!(producer.records().is_empty()); + backend.upload_offset(&token).await.unwrap(); + assert_eq!(producer.records().len(), 1); - // Completion emits exactly one write for the published object. backend .put_chunk(&token, 1, 1, stream::single("b")) .await .unwrap(); let records = producer.records(); - assert_eq!(records.len(), 1); - assert_eq!(records[0].op_type, OpType::Write); + assert_eq!(records.len(), 3); + assert_eq!( + records.iter().map(|r| r.op_type).collect::>(), + [OpType::Write, OpType::Write, OpType::Delete] + ); + assert_eq!(records[0].record_id, records[2].record_id); + assert_ne!(records[0].record_id, records[1].record_id); + assert_eq!( + records[1].size, + Some((json_len(&Metadata::default()) + 2) as u64) + ); + + producer.clear(); + let first = create_session(&backend, &id, 2).await; + let second = create_session(&backend, &id, 3).await; + backend.cancel_upload(&first).await.unwrap(); + assert!(backend.cancel_upload(&first).await.is_err()); + assert_eq!( + backend.upload_offset(&second).await.unwrap(), + UploadProgress::Incomplete { offset: 0 } + ); + backend.cancel_upload(&second).await.unwrap(); + let records = producer.records(); + assert_eq!(records.len(), 4); + assert_ne!(records[0].record_id, records[1].record_id); + assert_eq!(records[0].record_id, records[2].record_id); + assert_eq!(records[1].record_id, records[3].record_id); + assert_eq!(records[2].op_type, OpType::Delete); + assert_eq!(records[3].op_type, OpType::Delete); } #[tokio::test] diff --git a/objectstore-service/src/backend/local_fs.rs b/objectstore-service/src/backend/local_fs.rs index 751c64b1..4c6f0935 100644 --- a/objectstore-service/src/backend/local_fs.rs +++ b/objectstore-service/src/backend/local_fs.rs @@ -42,7 +42,8 @@ use crate::backend::common::{ SetExpiryResponse, }; use crate::change_stream::{ - ChangeStream, ChangeStreamFactory, CostTrackerStreamConfig, flush_change_stream, + ChangeStream, ChangeStreamFactory, ChangeTarget, CostTrackerStreamConfig, UPLOAD_SESSION_TTL, + flush_change_stream, }; use crate::error::{Error, ErrorKind, Result, ResultExt as _}; use crate::id::ObjectId; @@ -192,7 +193,7 @@ impl Backend for LocalFsBackend { draft.publish().await?; self.change_stream - .write(id, stored_size, metadata.time_expires); + .write(id.into(), stored_size, metadata.time_expires); Ok(()) } @@ -283,7 +284,7 @@ impl Backend for LocalFsBackend { draft.prepare().await?; draft.publish().await?; - self.change_stream.update(id, Some(expire_at)); + self.change_stream.update(id.into(), Some(expire_at)); Ok(SetExpiryResponse::Satisfied(expire_at)) } @@ -299,7 +300,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.delete(id.into()), Err(error) if error.kind() == io::ErrorKind::NotFound => { objectstore_log::debug!("Object not found"); } @@ -316,7 +317,7 @@ impl Backend for LocalFsBackend { &self, id: &ObjectId, metadata: &Metadata, - _upload_length: NonZeroU64, + upload_length: NonZeroU64, ) -> Result> { let upload_id = uuid::Uuid::now_v7(); let path = self.upload_path(upload_id); @@ -325,7 +326,16 @@ impl Backend for LocalFsBackend { "creating local-fs object directory", )?; UploadFile::create(&path, metadata).await?; - Ok(Some(upload_id.to_string())) + let token = upload_id.to_string(); + self.change_stream.write( + ChangeTarget::UploadSession { + object_id: id, + session_id: &token, + }, + upload_length.get(), + Some(Timestamp::now() + UPLOAD_SESSION_TTL), + ); + Ok(Some(token)) } // In this backend, if all the bytes of a resumable upload have been written but publication failed, @@ -377,7 +387,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); + .write((&session.object_id).into(), stored_size, expires_at); + self.change_stream.delete(session.into()); Ok(UploadProgress::Complete) } @@ -394,7 +405,10 @@ impl Backend for LocalFsBackend { let _guard = self.locks.acquire(&session.object_id).await?; let path = self.upload_path(upload_id); match tokio::fs::remove_file(&path).await { - Ok(()) => Ok(()), + Ok(()) => { + self.change_stream.delete(session.into()); + Ok(()) + } Err(error) if error.kind() == io::ErrorKind::NotFound => { Err(ErrorKind::UnknownUploadSession.into()) } @@ -741,7 +755,7 @@ impl MultipartUploadBackend for LocalFsBackend { drop(guard); self.change_stream - .write(id, stored_size, metadata.time_expires); + .write(id.into(), stored_size, metadata.time_expires); // Clean up multipart state tokio::fs::remove_dir_all(dir).await.context( @@ -1232,6 +1246,9 @@ mod tests { #[tokio::test] async fn resumable_publication_can_be_retried() { for query_offset in [false, true] { + #[cfg(feature = "storage-cogs")] + let (_tempdir, backend, producer) = make_backend_with_change_stream(); + #[cfg(not(feature = "storage-cogs"))] let (_tempdir, backend) = make_backend(); let id = make_id(); let token = upload_token(&backend, &id, 4).await; @@ -1244,6 +1261,12 @@ mod tests { .await .unwrap_err(); assert_eq!(error.kind(), ErrorKind::BackendFailure); + #[cfg(feature = "storage-cogs")] + assert_eq!( + producer.records().len(), + 1, + "failed publication leaves the session record" + ); tokio::fs::remove_dir(&object_path).await.unwrap(); let progress = if query_offset { @@ -1252,6 +1275,15 @@ mod tests { backend.put_chunk(&token, 4, 0, stream::single("")).await }; assert_eq!(progress.unwrap(), UploadProgress::Complete); + #[cfg(feature = "storage-cogs")] + { + let records = producer.records(); + assert_eq!( + records.iter().map(|r| r.op_type).collect::>(), + [OpType::Write, OpType::Write, OpType::Delete] + ); + assert_eq!(records[0].record_id, records[2].record_id); + } let (_, _, payload) = backend .get_object(&id, Timestamp::now(), None) .await @@ -2264,6 +2296,7 @@ mod tests { }; let payload = b"oh hai!"; let upload_length = NonZeroU64::new(payload.len() as u64).unwrap(); + let before = Timestamp::now(); let token = backend .create_upload_session(&id, &metadata, upload_length) .await @@ -2275,17 +2308,30 @@ mod tests { backend_token: token, }; - // Session creation alone does not report a stored object. - assert!(producer.records().is_empty()); + let created = producer.records(); + assert_eq!(created.len(), 1); + assert_eq!(created[0].size, Some(payload.len() as u64)); + let expiration = created[0].expiration_time.unwrap() as u64; + assert!( + ((before + UPLOAD_SESSION_TTL).as_micros() + ..=(Timestamp::now() + UPLOAD_SESSION_TTL).as_micros()) + .contains(&expiration) + ); + backend + .put_chunk(&token, 0, 1, stream::single(payload[..1].to_vec())) + .await + .unwrap(); + backend.upload_offset(&token).await.unwrap(); + assert_eq!(producer.records().len(), 1); // Completing the upload reports the published file's full stored size. assert_eq!( backend .put_chunk( &token, - 0, - payload.len() as u64, - stream::single(payload.to_vec()), + 1, + payload.len() as u64 - 1, + stream::single(payload[1..].to_vec()), ) .await .unwrap(), @@ -2296,10 +2342,27 @@ mod tests { .await .unwrap(); let records = producer.records(); - assert_eq!(records.len(), 1); - assert_eq!(records[0].op_type, OpType::Write); - assert_eq!(records[0].size, Some(file.len() as u64)); - assert!(records[0].expiration_time.is_some()); + assert_eq!(records.len(), 3); + assert_eq!( + records.iter().map(|r| r.op_type).collect::>(), + [OpType::Write, OpType::Write, OpType::Delete] + ); + assert_eq!(records[0].record_id, records[2].record_id); + assert_ne!(records[0].record_id, records[1].record_id); + assert_eq!(records[1].size, Some(file.len() as u64)); + assert_eq!( + records[1].expiration_time, + metadata.time_expires.map(|t| t.as_micros() as i64) + ); + + producer.clear(); + let canceled = upload_token(&backend, &id, 10).await; + backend.cancel_upload(&canceled).await.unwrap(); + assert!(backend.cancel_upload(&canceled).await.is_err()); + let records = producer.records(); + assert_eq!(records.len(), 2); + assert_eq!(records[1].op_type, OpType::Delete); + assert_eq!(records[0].record_id, records[1].record_id); } #[cfg(feature = "storage-cogs")] diff --git a/objectstore-service/src/backend/s3_compatible.rs b/objectstore-service/src/backend/s3_compatible.rs index b0bcac04..ae5ca961 100644 --- a/objectstore-service/src/backend/s3_compatible.rs +++ b/objectstore-service/src/backend/s3_compatible.rs @@ -398,7 +398,7 @@ impl Backend for S3CompatibleBackend { .await; self.change_stream.write( - id, + id.into(), metadata_size + payload_size.load(Ordering::Relaxed), metadata.time_expires, ); @@ -481,7 +481,7 @@ impl Backend for S3CompatibleBackend { .update_metadata(id, &metadata, expire_at, &etag) .await?; if matches!(outcome, SetExpiryResponse::Satisfied(_)) { - self.change_stream.update(id, Some(expire_at)); + self.change_stream.update(id.into(), Some(expire_at)); } Ok(outcome) @@ -516,7 +516,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.delete(id.into()); Ok(()) } diff --git a/objectstore-service/src/backend/tiered.rs b/objectstore-service/src/backend/tiered.rs index 97bee063..f458fcd8 100644 --- a/objectstore-service/src/backend/tiered.rs +++ b/objectstore-service/src/backend/tiered.rs @@ -103,8 +103,8 @@ //! written to the long-term backend. A resumable upload remains inaccessible until //! the long-term backend completes and a high-volume tombstone is committed. //! -//! Each upload has a separate HV marker with a fixed five-day lifetime that serves -//! as the source of truth for whether the upload logically exists. +//! Each upload has a separate HV marker with a lifetime set by `RESUMABLE_UPLOAD_TTL`. +//! It serves as the source of truth for whether the upload logically exists. //! This marker is consumed upon cancellation or the first submission of a final chunk, making //! finalization one-shot. @@ -1256,6 +1256,7 @@ mod tests { use crate::backend::gcs::{GcsBackend, GcsConfig}; use crate::backend::in_memory::InMemoryBackend; use crate::backend::testing::{Hooks, TestBackend}; + #[cfg(not(feature = "storage-cogs"))] use crate::change_stream::ChangeStreamFactory; use crate::error::Error; use crate::id::ObjectContext; @@ -1396,11 +1397,20 @@ mod tests { #[tokio::test] async fn resumable_bigtable_and_gcs() -> anyhow::Result<()> { + #[cfg(feature = "storage-cogs")] + let (streams, producer) = crate::change_stream::dummy_factory(); + #[cfg(not(feature = "storage-cogs"))] let streams = ChangeStreamFactory::default(); let lt = GcsBackend::new( GcsConfig { endpoint: Some("http://localhost:8087".into()), bucket: "test-bucket".into(), + #[cfg(feature = "storage-cogs")] + cogs: Some(crate::change_stream::CostTrackerStreamConfig { + shared_resource_id: "gcs_objectstore".into(), + sample_rate: 1.0, + }), + #[cfg(not(feature = "storage-cogs"))] cogs: None, }, &streams, @@ -1414,6 +1424,12 @@ mod tests { table_name: "objectstore".into(), connections: None, rpc_timeout: Duration::from_secs(2), + #[cfg(feature = "storage-cogs")] + cogs: Some(crate::change_stream::CostTrackerStreamConfig { + shared_resource_id: "bigtable_objectstore".into(), + sample_rate: 1.0, + }), + #[cfg(not(feature = "storage-cogs"))] cogs: None, }, &streams, @@ -1422,6 +1438,8 @@ mod tests { let storage = TieredStorage::new(Box::new(hv), Box::new(lt), Box::new(NoopChangeLog)); let id = make_id(&format!("tiered-resumable-{}", uuid::Uuid::now_v7())); let payload = vec![b'a'; BACKEND_SIZE_THRESHOLD + 1]; + #[cfg(feature = "storage-cogs")] + let before = Timestamp::now(); let token = resumable_token(&storage, &id, &Metadata::default(), payload.len() as u64).await; @@ -1438,6 +1456,22 @@ mod tests { ); let revision = upload_revision(&token); + #[cfg(feature = "storage-cogs")] + { + let records = producer.records(); + assert_eq!(records.len(), 2); + assert_eq!(records[0].size, Some(payload.len() as u64)); + assert_eq!( + records[1].size, + Some(revision.as_upload_path().to_string().len() as u64 + 1) + ); + let expiry = records[1].expiration_time.unwrap() as u64; + assert!( + ((before + RESUMABLE_UPLOAD_TTL).as_micros() + ..=(Timestamp::now() + RESUMABLE_UPLOAD_TTL).as_micros()) + .contains(&expiry) + ); + } assert!( storage .inner @@ -1478,6 +1512,13 @@ mod tests { ); assert!(storage.get_metadata(&id, Timestamp::now()).await?.is_none()); + #[cfg(feature = "storage-cogs")] + assert_eq!( + producer.records().len(), + 2, + "chunks and queries do not report changes" + ); + // A completed upload creates a logical object. assert_eq!( storage @@ -1506,6 +1547,45 @@ mod tests { .await? .unwrap(); assert_eq!(stream::read_to_vec(body).await?, payload); + #[cfg(feature = "storage-cogs")] + { + use objectstore_inventory_tracker::OpType::{Delete, Write}; + let records = producer.records(); + assert_eq!( + records + .iter() + .map(|r| (r.shared_resource_id.as_str(), r.op_type)) + .collect::>(), + [ + ("gcs_objectstore", Write), + ("bigtable_objectstore", Write), + ("bigtable_objectstore", Delete), + ("gcs_objectstore", Write), + ("gcs_objectstore", Delete), + ("bigtable_objectstore", Write), + ] + ); + assert_eq!(records[0].record_id, records[4].record_id); + assert_eq!(records[1].record_id, records[2].record_id); + assert_ne!(records[0].record_id, records[3].record_id); + assert_ne!(records[1].record_id, records[5].record_id); + + producer.clear(); + let canceled = + resumable_token(&storage, &id, &Metadata::default(), payload.len() as u64).await; + storage.cancel_upload(&canceled).await?; + assert_eq!( + storage.cancel_upload(&canceled).await.unwrap_err().kind(), + ErrorKind::UploadSessionGone + ); + let records = producer.records(); + assert_eq!( + records.iter().map(|r| r.op_type).collect::>(), + [Write, Write, Delete, Delete] + ); + assert_eq!(records[0].record_id, records[3].record_id); + assert_eq!(records[1].record_id, records[2].record_id); + } Ok(()) } diff --git a/objectstore-service/src/change_stream/cost_tracker.rs b/objectstore-service/src/change_stream/cost_tracker.rs index 33d805ce..11bc7af5 100644 --- a/objectstore-service/src/change_stream/cost_tracker.rs +++ b/objectstore-service/src/change_stream/cost_tracker.rs @@ -6,15 +6,24 @@ use std::time::{Duration, SystemTime}; use objectstore_inventory_tracker::{BoxError, InventoryTracker, Producer}; use crate::change_stream::{ - ChangeStream, CostTrackerStreamConfig, SCOPE_ORGANIZATION, SCOPE_PROJECT, scope_id, + ChangeStream, ChangeTarget, CostTrackerStreamConfig, SCOPE_ORGANIZATION, SCOPE_PROJECT, + scope_id, }; use objectstore_types::time::Timestamp; use crate::id::ObjectId; -/// Reports through an [`InventoryTracker`], which hashes each [`ObjectId`] both to -/// anonymize it and to decide whether it is sampled. See [`objectstore_inventory_tracker`] -/// for the record format. +/// Reports through an [`InventoryTracker`], which hashes each target identity both to +/// anonymize it and to decide whether it is sampled. Session identities include the object +/// path and session ID, keeping their records and sampling separate from published objects. +/// See [`objectstore_inventory_tracker`] for the record format. +/// +/// Object identities remain their storage paths. Session identities are JSON tuples of +/// `("upload_session", object_storage_path, session_id)`. Every operation on a session +/// uses the same identity. Sampling decisions are consistent while the backend +/// [`sample_rate`](CostTrackerStreamConfig::sample_rate) is unchanged; configuration changes +/// can change the decision between a session write and delete. Object and session sampling +/// decisions may differ. Neither raw paths nor session tokens are emitted. /// /// Logs, counts, and swallows errors returned by the [`InventoryTracker`]. pub struct CostTrackerStream { @@ -67,15 +76,39 @@ impl fmt::Debug for CostTrackerStream

{ } } +impl<'a> ChangeTarget<'a> { + /// Returns the attribution ID and stable input to inventory hashing and sampling. + fn inventory_identity(self) -> (&'a ObjectId, String) { + match self { + Self::Object(id) => (id, id.as_storage_path().to_string()), + Self::UploadSession { + object_id, + session_id, + } => { + // This tuple is the permanent session identity contract, just as the + // storage path is for objects. Neither the path nor token is emitted. + let key = serde_json::to_string(&( + "upload_session", + object_id.as_storage_path(), + session_id, + )) + .expect("session identity is serializable"); + (object_id, key) + } + } + } +} + #[async_trait::async_trait] impl

ChangeStream for CostTrackerStream

where P: Producer + Clone + Send + Sync + 'static, P::Error: Into + Send + 'static, { - fn write(&self, id: &ObjectId, size: u64, expires_at: Option) { + fn write(&self, target: ChangeTarget<'_>, size: u64, expires_at: Option) { + let (id, key) = target.inventory_identity(); let result = self.tracker.write( - &id.as_storage_path().to_string(), + &key, id.usecase(), size, SystemTime::now(), @@ -86,9 +119,10 @@ where self.swallow("write", result); } - fn update(&self, id: &ObjectId, expires_at: Option) { + fn update(&self, target: ChangeTarget<'_>, expires_at: Option) { + let (id, key) = target.inventory_identity(); let result = self.tracker.update( - &id.as_storage_path().to_string(), + &key, id.usecase(), SystemTime::now(), expires_at.map(Into::into), @@ -98,12 +132,9 @@ where self.swallow("update", result); } - fn delete(&self, id: &ObjectId) { - let result = self.tracker.delete( - &id.as_storage_path().to_string(), - id.usecase(), - SystemTime::now(), - ); + fn delete(&self, target: ChangeTarget<'_>) { + let (id, key) = target.inventory_identity(); + let result = self.tracker.delete(&key, id.usecase(), SystemTime::now()); self.swallow("delete", result); } @@ -140,7 +171,7 @@ mod tests { let (producer, stream) = stream(1.0); let id = object_id("attachments/org.17/project.42/objects/abc"); - stream.write(&id, 4096, None); + stream.write((&id).into(), 4096, None); let record = &producer.records()[0]; assert_eq!(record.shared_resource_id, "bigtable_objectstore"); @@ -156,7 +187,7 @@ mod tests { let (producer, stream) = stream(1.0); let id = object_id("attachments/org.17/project.42/objects/abc"); - stream.write(&id, 4096, None); + stream.write((&id).into(), 4096, None); let record_id = &producer.records()[0].record_id; assert_ne!(record_id, &id.as_storage_path().to_string()); @@ -172,7 +203,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.write((&object_id(path)).into(), 1, None); } let records = producer.records(); @@ -205,9 +236,9 @@ mod tests { 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.write((&id).into(), 10, None); + stream.update((&id).into(), Some(Timestamp::now())); + stream.delete((&id).into()); let records = producer.records(); assert_eq!(records.len(), 3); @@ -220,12 +251,12 @@ mod tests { let (producer, stream) = stream(1.0); stream.write( - &object_id("attachments/org.1/project.2/objects/abc/0199aaaa"), + (&object_id("attachments/org.1/project.2/objects/abc/0199aaaa")).into(), 1, None, ); stream.write( - &object_id("attachments/org.1/project.2/objects/abc/0199bbbb"), + (&object_id("attachments/org.1/project.2/objects/abc/0199bbbb")).into(), 1, None, ); @@ -240,8 +271,8 @@ mod tests { 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.update((&id).into(), Some(expires)); + stream.delete((&id).into()); let records = producer.records(); assert_eq!(records[0].op_type, OpType::Update); @@ -252,15 +283,51 @@ mod tests { assert_eq!(records[1].expiration_time, None); } + #[test] + fn session_sampling_is_consistent_and_separate_from_objects() { + let (producer, stream) = stream(0.5); + let id = object_id("attachments/org.17/project.42/objects/abc"); + let mut sampled = 0; + let mut record_ids = std::collections::HashSet::new(); + for i in 0..32 { + let session_id = format!("session-{i}"); + let target = ChangeTarget::UploadSession { + object_id: &id, + session_id: &session_id, + }; + stream.write(target, 10, None); + stream.delete(target); + let records = producer.records(); + if !records.is_empty() { + sampled += 1; + assert_eq!(records.len(), 2); + assert_eq!(records[0].op_type, OpType::Write); + assert_eq!(records[1].op_type, OpType::Delete); + assert_eq!(records[0].record_id, records[1].record_id); + assert!(record_ids.insert(records[0].record_id.clone())); + assert_eq!(records[0].organization_id, Some(17)); + assert_eq!(records[0].project_id, Some(42)); + assert_eq!(records[0].app_feature, "attachments"); + assert_eq!(records[0].sample_rate, 0.5); + } + producer.clear(); + } + assert!(sampled > 0 && sampled < 32); + + let (producer, stream) = self::stream(1.0); + stream.write((&id).into(), 10, None); + assert!(!record_ids.contains(&producer.records()[0].record_id)); + } + #[test] 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.write((&id).into(), 1, None); + stream.update((&id).into(), None); + stream.delete((&id).into()); } 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..5a69c610 100644 --- a/objectstore-service/src/change_stream/factory.rs +++ b/objectstore-service/src/change_stream/factory.rs @@ -156,7 +156,9 @@ mod tests { assert!(reports(&stream)); - stream.delete(&crate::id::ObjectId::from_storage_path("attachments/objects/abc").unwrap()); + stream.delete( + (&crate::id::ObjectId::from_storage_path("attachments/objects/abc").unwrap()).into(), + ); 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..9747f427 100644 --- a/objectstore-service/src/change_stream/mod.rs +++ b/objectstore-service/src/change_stream/mod.rs @@ -4,6 +4,21 @@ //! the service describes where those records go with a [`CostTrackerConfig`], shared by //! every backend. [`ChangeStreamFactory`] pairs the two into a [`ChangeStream`]. //! +//! # Resumable uploads +//! +//! Upload backends report a session write after creation with exactly the advertised upload +//! length and an accounting expiration set by `UPLOAD_SESSION_TTL`, independent of object +//! expiration. Partial chunks and incomplete offset queries emit nothing. Successful publication +//! reports the object's actual stored size and expiration, followed by a session delete. Completion +//! discovered through an offset query follows the same order. Successful cancellation also +//! reports a session delete. The accounting deadline does not change backend cleanup behavior. +//! +//! High-volume markers report their stored size and actual deadline, set by +//! `backend::tiered::RESUMABLE_UPLOAD_TTL`. Only a successful +//! conditional marker deletion emits a delete. Tiered deletes its marker to claim the upload +//! before finalization, so that marker delete precedes the long-term object write and session +//! delete. Failed operations leave accounting intact until successful cleanup or expiration. +//! //! Behind the `storage-cogs` feature. Without it every backend gets a [`NoopStream`] and //! the transport is left out of the binary. @@ -16,6 +31,7 @@ use serde::{Deserialize, Serialize}; use objectstore_types::time::Timestamp; use crate::id::ObjectId; +use crate::resumable::Session; #[cfg(feature = "storage-cogs")] mod cost_tracker; @@ -33,6 +49,52 @@ pub(crate) use factory::dummy_factory; /// How long a backend waits for reported records to be handed off during shutdown. pub const FLUSH_TIMEOUT: Duration = Duration::from_secs(2); +/// Accounting lifetime for an upload's advertised size, independent of backend cleanup. +pub(crate) const UPLOAD_SESSION_TTL: Duration = Duration::from_hours(7 * 24); + +/// The object or upload session whose storage changed. +/// +/// Session identities are separate from objects, including concurrent uploads to the same object. +/// Upload backends use their opaque backend token as `session_id`; high-volume markers use +/// [`ChangeTarget::upload_marker`]. +#[derive(Clone, Copy, Debug)] +pub enum ChangeTarget<'a> { + /// A published object, including a high-volume redirect tombstone. + Object(&'a ObjectId), + /// An in-progress upload or its high-volume marker. + UploadSession { + /// Object identity supplying the usecase and scopes. + object_id: &'a ObjectId, + /// Stable identity for this particular upload. + session_id: &'a str, + }, +} + +impl<'a> ChangeTarget<'a> { + /// Identifies the high-volume marker for an upload's unique revision. + pub fn upload_marker(revision: &'a ObjectId) -> Self { + Self::UploadSession { + object_id: revision, + session_id: &revision.key, + } + } +} + +impl<'a> From<&'a ObjectId> for ChangeTarget<'a> { + fn from(id: &'a ObjectId) -> Self { + Self::Object(id) + } +} + +impl<'a> From<&'a Session> for ChangeTarget<'a> { + fn from(session: &'a Session) -> Self { + Self::UploadSession { + object_id: &session.object_id, + session_id: &session.backend_token, + } + } +} + /// Scope key holding the Sentry organization ID. #[cfg(feature = "storage-cogs")] const SCOPE_ORGANIZATION: &str = "org"; @@ -75,19 +137,23 @@ fn default_sample_rate() -> f64 { 1.0 } -/// Publishes the changes a single backend makes to the objects it stores. +/// Publishes the changes a single backend makes to its objects and upload sessions. /// /// See [module docs](self). #[async_trait::async_trait] pub trait ChangeStream: fmt::Debug + Send + Sync + 'static { - /// Reports that `id` now occupies `size` bytes. Used for new writes and overwrites. - fn write(&self, id: &ObjectId, size: u64, expires_at: Option); + /// Reports that `target` now occupies `size` bytes. Used for new writes and overwrites. + /// + /// Upload sessions optimistically report the advertised upload length with an accounting + /// expiration set by `UPLOAD_SESSION_TTL`. Markers report their stored size and actual + /// expiration instead. + fn write(&self, target: ChangeTarget<'_>, size: u64, expires_at: Option); - /// Reports that `id`'s expiration moved, with its stored size unchanged. - fn update(&self, id: &ObjectId, expires_at: Option); + /// Reports that `target`'s expiration moved, with its stored size unchanged. + fn update(&self, target: ChangeTarget<'_>, expires_at: Option); - /// Reports that `id` was deleted explicitly. Does not account for automatic GC. - fn delete(&self, id: &ObjectId); + /// Reports that `target` was deleted explicitly. Does not account for automatic GC. + fn delete(&self, target: ChangeTarget<'_>); /// Blocks until reported records have been delivered, or `timeout` elapses. /// @@ -113,11 +179,11 @@ pub struct NoopStream; #[async_trait::async_trait] impl ChangeStream for NoopStream { - fn write(&self, _id: &ObjectId, _size: u64, _expires_at: Option) {} + fn write(&self, _target: ChangeTarget<'_>, _size: u64, _expires_at: Option) {} - fn update(&self, _id: &ObjectId, _expires_at: Option) {} + fn update(&self, _target: ChangeTarget<'_>, _expires_at: Option) {} - fn delete(&self, _id: &ObjectId) {} + fn delete(&self, _target: ChangeTarget<'_>) {} async fn join(&self, _timeout: Duration) {} } diff --git a/objectstore-service/src/id.rs b/objectstore-service/src/id.rs index 4d8459f0..2b6104c8 100644 --- a/objectstore-service/src/id.rs +++ b/objectstore-service/src/id.rs @@ -237,6 +237,15 @@ pub struct AsStoragePath<'a> { namespace: &'static str, } +impl serde::Serialize for AsStoragePath<'_> { + fn serialize(&self, serializer: S) -> Result + where + S: serde::Serializer, + { + serializer.collect_str(self) + } +} + impl fmt::Display for AsStoragePath<'_> { fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { write!(f, "{}/", self.inner.context.usecase)?; @@ -266,6 +275,10 @@ mod tests { let path = object_id.as_storage_path().to_string(); assert_eq!(path, "testing/org.12345/project.1337/objects/foo/bar"); + assert_eq!( + serde_json::to_string(&object_id.as_storage_path()).unwrap(), + serde_json::to_string(&path).unwrap(), + ); } #[test] From 0294b0d3c5e42802f31e41bf550f820a17fc2c9a Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Thu, 8 Oct 2026 15:31:48 +0200 Subject: [PATCH 02/11] fix(cogs): Exclude incomplete GCS upload sessions Default upload-session accounting to enabled in CostTrackerStream and disable it for GCS at construction. Preserve generic session lifecycle events and existing inventory identities. --- objectstore-service/docs/architecture.md | 4 +- objectstore-service/src/backend/gcs.rs | 46 +++++-------------- objectstore-service/src/backend/tiered.rs | 20 +++----- .../src/change_stream/cost_tracker.rs | 45 ++++++++++++++++++ .../src/change_stream/factory.rs | 23 ++++++++-- objectstore-service/src/change_stream/mod.rs | 3 ++ 6 files changed, 90 insertions(+), 51 deletions(-) diff --git a/objectstore-service/docs/architecture.md b/objectstore-service/docs/architecture.md index 2d7e5d9d..d3348871 100644 --- a/objectstore-service/docs/architecture.md +++ b/objectstore-service/docs/architecture.md @@ -166,7 +166,9 @@ row for the actual object and its size in long-term backend, and a separate row for the tombstone and the tombstone's size in the high-volume backend. While a resumable upload is in progress, separate session rows account for its advertised size in the upload backend and its marker's stored size in high-volume -storage. See the [change stream module](change_stream) for their lifetimes. +storage. GCS excludes incomplete uploads in its in-process cost-tracking adapter, +since they do not incur storage charges; the generic stream still reports their +lifecycle. See the [change stream module](change_stream) for their lifetimes. Because the change stream does not observe automatic garbage collection, expired records must be filtered out when querying the inventory table. diff --git a/objectstore-service/src/backend/gcs.rs b/objectstore-service/src/backend/gcs.rs index 5bc67633..0375809e 100644 --- a/objectstore-service/src/backend/gcs.rs +++ b/objectstore-service/src/backend/gcs.rs @@ -518,7 +518,8 @@ impl GcsBackend { bucket, cogs, } = config; - let change_stream = streams.build(cogs.as_ref()); + // Incomplete GCS uploads do not contribute to storage COGS. + let change_stream = streams.build_with_upload_sessions(cogs.as_ref(), false); let token_provider = if endpoint.is_none() { Some(PrefetchingTokenProvider::gcp_auth(TOKEN_SCOPES).await?) @@ -3259,22 +3260,13 @@ mod tests { time_expires: Some(Timestamp::now() + Duration::from_secs(3600)), ..Default::default() }; - let before = Timestamp::now(); let token = backend .create_upload_session(&id, &metadata, nonzero(payload.len() as u64)) .await?; - let created = producer.records(); - assert_eq!(created.len(), 1); - assert_eq!(created[0].size, Some(payload.len() as u64)); - let expiration = created[0].expiration_time.unwrap() as u64; - assert!( - ((before + UPLOAD_SESSION_TTL).as_micros() - ..=(Timestamp::now() + UPLOAD_SESSION_TTL).as_micros()) - .contains(&expiration) - ); + assert!(producer.records().is_empty()); backend.upload_offset(&token).await?; - assert_eq!(producer.records().len(), 1); + assert!(producer.records().is_empty()); assert_eq!( backend @@ -3289,19 +3281,14 @@ mod tests { ); let records = producer.records(); - assert_eq!(records.len(), 3); - assert_eq!( - records.iter().map(|r| r.op_type).collect::>(), - [OpType::Write, OpType::Write, OpType::Delete] - ); - assert_eq!(records[0].record_id, records[2].record_id); - assert_ne!(records[0].record_id, records[1].record_id); + assert_eq!(records.len(), 1); + assert_eq!(records[0].op_type, OpType::Write); assert_eq!( - records[1].size, + records[0].size, Some(payload.len() as u64 + GcsObject::from_metadata(&metadata).metadata_size()) ); assert_eq!( - records[1].expiration_time, + records[0].expiration_time, metadata.time_expires.map(|t| t.as_micros() as i64) ); producer.clear(); @@ -3310,11 +3297,7 @@ mod tests { .await?; backend.cancel_upload(&canceled).await?; backend.cancel_upload(&canceled).await?; - let records = producer.records(); - assert_eq!(records.len(), 3); - assert_eq!(records[1].op_type, OpType::Delete); - assert_eq!(records[2].op_type, OpType::Delete); - assert!(records.iter().all(|r| r.record_id == records[0].record_id)); + assert!(producer.records().is_empty()); Ok(()) } @@ -3352,15 +3335,10 @@ mod tests { ); let records = producer.records(); - assert_eq!(records.len(), 3); - assert_eq!( - records.iter().map(|r| r.op_type).collect::>(), - [OpType::Write, OpType::Write, OpType::Delete] - ); - assert_eq!(records[0].record_id, records[2].record_id); - assert_ne!(records[0].record_id, records[1].record_id); + assert_eq!(records.len(), 1); + assert_eq!(records[0].op_type, OpType::Write); assert_eq!( - records[1].size, + records[0].size, Some(payload.len() as u64 + GcsObject::from_metadata(&metadata).metadata_size()) ); Ok(()) diff --git a/objectstore-service/src/backend/tiered.rs b/objectstore-service/src/backend/tiered.rs index f458fcd8..6d8dc19a 100644 --- a/objectstore-service/src/backend/tiered.rs +++ b/objectstore-service/src/backend/tiered.rs @@ -1459,13 +1459,12 @@ mod tests { #[cfg(feature = "storage-cogs")] { let records = producer.records(); - assert_eq!(records.len(), 2); - assert_eq!(records[0].size, Some(payload.len() as u64)); + assert_eq!(records.len(), 1); assert_eq!( - records[1].size, + records[0].size, Some(revision.as_upload_path().to_string().len() as u64 + 1) ); - let expiry = records[1].expiration_time.unwrap() as u64; + let expiry = records[0].expiration_time.unwrap() as u64; assert!( ((before + RESUMABLE_UPLOAD_TTL).as_micros() ..=(Timestamp::now() + RESUMABLE_UPLOAD_TTL).as_micros()) @@ -1515,7 +1514,7 @@ mod tests { #[cfg(feature = "storage-cogs")] assert_eq!( producer.records().len(), - 2, + 1, "chunks and queries do not report changes" ); @@ -1557,18 +1556,14 @@ mod tests { .map(|r| (r.shared_resource_id.as_str(), r.op_type)) .collect::>(), [ - ("gcs_objectstore", Write), ("bigtable_objectstore", Write), ("bigtable_objectstore", Delete), ("gcs_objectstore", Write), - ("gcs_objectstore", Delete), ("bigtable_objectstore", Write), ] ); - assert_eq!(records[0].record_id, records[4].record_id); - assert_eq!(records[1].record_id, records[2].record_id); + assert_eq!(records[0].record_id, records[1].record_id); assert_ne!(records[0].record_id, records[3].record_id); - assert_ne!(records[1].record_id, records[5].record_id); producer.clear(); let canceled = @@ -1581,10 +1576,9 @@ mod tests { let records = producer.records(); assert_eq!( records.iter().map(|r| r.op_type).collect::>(), - [Write, Write, Delete, Delete] + [Write, Delete] ); - assert_eq!(records[0].record_id, records[3].record_id); - assert_eq!(records[1].record_id, records[2].record_id); + assert_eq!(records[0].record_id, records[1].record_id); } Ok(()) } diff --git a/objectstore-service/src/change_stream/cost_tracker.rs b/objectstore-service/src/change_stream/cost_tracker.rs index 11bc7af5..fc9fccc5 100644 --- a/objectstore-service/src/change_stream/cost_tracker.rs +++ b/objectstore-service/src/change_stream/cost_tracker.rs @@ -25,9 +25,14 @@ use crate::id::ObjectId; /// can change the decision between a session write and delete. Object and session sampling /// decisions may differ. Neither raw paths nor session tokens are emitted. /// +/// Upload sessions are tracked by default. GCS disables their accounting at stream +/// construction because incomplete uploads do not incur storage charges. Generic +/// session lifecycle events remain available to other change stream implementations. +/// /// Logs, counts, and swallows errors returned by the [`InventoryTracker`]. pub struct CostTrackerStream { tracker: InventoryTracker

, + pub(super) track_upload_sessions: bool, } impl CostTrackerStream

{ @@ -39,6 +44,7 @@ impl CostTrackerStream

{ &config.shared_resource_id, config.sample_rate, ), + track_upload_sessions: true, } } @@ -72,6 +78,7 @@ impl fmt::Debug for CostTrackerStream

{ f.debug_struct("CostTrackerStream") .field("shared_resource_id", &self.tracker.shared_resource_id()) .field("sample_rate", &self.tracker.sample_rate()) + .field("track_upload_sessions", &self.track_upload_sessions) .finish() } } @@ -106,6 +113,9 @@ where P::Error: Into + Send + 'static, { fn write(&self, target: ChangeTarget<'_>, size: u64, expires_at: Option) { + if !self.track_upload_sessions && matches!(target, ChangeTarget::UploadSession { .. }) { + return; + } let (id, key) = target.inventory_identity(); let result = self.tracker.write( &key, @@ -120,6 +130,9 @@ where } fn update(&self, target: ChangeTarget<'_>, expires_at: Option) { + if !self.track_upload_sessions && matches!(target, ChangeTarget::UploadSession { .. }) { + return; + } let (id, key) = target.inventory_identity(); let result = self.tracker.update( &key, @@ -133,6 +146,9 @@ where } fn delete(&self, target: ChangeTarget<'_>) { + if !self.track_upload_sessions && matches!(target, ChangeTarget::UploadSession { .. }) { + return; + } let (id, key) = target.inventory_identity(); let result = self.tracker.delete(&key, id.usecase(), SystemTime::now()); self.swallow("delete", result); @@ -319,6 +335,35 @@ mod tests { assert!(!record_ids.contains(&producer.records()[0].record_id)); } + #[test] + fn upload_sessions_can_be_excluded_without_affecting_objects() { + let (producer, mut stream) = stream(1.0); + let id = object_id("attachments/org.17/project.42/objects/abc"); + let session = ChangeTarget::UploadSession { + object_id: &id, + session_id: "upload-token", + }; + stream.track_upload_sessions = false; + stream.write(session, 10, None); + stream.update(session, None); + stream.delete(session); + assert!(producer.records().is_empty()); + + stream.write((&id).into(), 10, None); + stream.update((&id).into(), None); + stream.delete((&id).into()); + let records = producer.records(); + assert_eq!(records.len(), 3); + assert_eq!(records[0].op_type, OpType::Write); + assert_eq!(records[1].op_type, OpType::Update); + assert_eq!(records[2].op_type, OpType::Delete); + assert!( + records + .iter() + .all(|record| record.record_id == records[0].record_id) + ); + } + #[test] fn a_listener_sampled_at_zero_reports_nothing() { let (producer, stream) = stream(0.0); diff --git a/objectstore-service/src/change_stream/factory.rs b/objectstore-service/src/change_stream/factory.rs index 5a69c610..2dced40c 100644 --- a/objectstore-service/src/change_stream/factory.rs +++ b/objectstore-service/src/change_stream/factory.rs @@ -48,10 +48,23 @@ impl ChangeStreamFactory { } /// Builds the stream `config` asks for, or a [`NoopStream`] if it cannot be built. - #[cfg(feature = "storage-cogs")] pub fn build(&self, config: Option<&CostTrackerStreamConfig>) -> Arc { + self.build_with_upload_sessions(config, true) + } + + /// Builds a stream with the backend's upload-session accounting policy. + #[cfg(feature = "storage-cogs")] + pub(crate) fn build_with_upload_sessions( + &self, + config: Option<&CostTrackerStreamConfig>, + track_upload_sessions: bool, + ) -> Arc { match (config, self.producer.clone()) { - (Some(config), Some(producer)) => Arc::new(CostTrackerStream::new(producer, config)), + (Some(config), Some(producer)) => { + let mut stream = CostTrackerStream::new(producer, config); + stream.track_upload_sessions = track_upload_sessions; + Arc::new(stream) + } (None, None) => Arc::new(NoopStream), (c, p) => { objectstore_log::warn!( @@ -66,7 +79,11 @@ impl ChangeStreamFactory { /// Reporting is not compiled in, so every backend reports nothing. #[cfg(not(feature = "storage-cogs"))] - pub fn build(&self, _config: Option<&CostTrackerStreamConfig>) -> Arc { + pub(crate) fn build_with_upload_sessions( + &self, + _config: Option<&CostTrackerStreamConfig>, + _track_upload_sessions: bool, + ) -> Arc { Arc::new(NoopStream) } } diff --git a/objectstore-service/src/change_stream/mod.rs b/objectstore-service/src/change_stream/mod.rs index 9747f427..3a88d629 100644 --- a/objectstore-service/src/change_stream/mod.rs +++ b/objectstore-service/src/change_stream/mod.rs @@ -4,6 +4,9 @@ //! the service describes where those records go with a [`CostTrackerConfig`], shared by //! every backend. [`ChangeStreamFactory`] pairs the two into a [`ChangeStream`]. //! +//! Cost tracking includes upload sessions by default. GCS disables session accounting in +//! the cost-tracking adapter while still publishing the generic lifecycle events. +//! //! # Resumable uploads //! //! Upload backends report a session write after creation with exactly the advertised upload From 34f572c3cdaa57a9a0a4afe2ed795b4a3de3469f Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Thu, 8 Oct 2026 15:51:48 +0200 Subject: [PATCH 03/11] revert(cogs): Restore GCS upload session records Leave storage-cost filtering to the downstream COGS consumer. This reverts commit 0294b0d3c5e42802f31e41bf550f820a17fc2c9a. --- objectstore-service/docs/architecture.md | 4 +- objectstore-service/src/backend/gcs.rs | 46 ++++++++++++++----- objectstore-service/src/backend/tiered.rs | 20 +++++--- .../src/change_stream/cost_tracker.rs | 45 ------------------ .../src/change_stream/factory.rs | 23 ++-------- objectstore-service/src/change_stream/mod.rs | 3 -- 6 files changed, 51 insertions(+), 90 deletions(-) diff --git a/objectstore-service/docs/architecture.md b/objectstore-service/docs/architecture.md index d3348871..2d7e5d9d 100644 --- a/objectstore-service/docs/architecture.md +++ b/objectstore-service/docs/architecture.md @@ -166,9 +166,7 @@ row for the actual object and its size in long-term backend, and a separate row for the tombstone and the tombstone's size in the high-volume backend. While a resumable upload is in progress, separate session rows account for its advertised size in the upload backend and its marker's stored size in high-volume -storage. GCS excludes incomplete uploads in its in-process cost-tracking adapter, -since they do not incur storage charges; the generic stream still reports their -lifecycle. See the [change stream module](change_stream) for their lifetimes. +storage. See the [change stream module](change_stream) for their lifetimes. Because the change stream does not observe automatic garbage collection, expired records must be filtered out when querying the inventory table. diff --git a/objectstore-service/src/backend/gcs.rs b/objectstore-service/src/backend/gcs.rs index 0375809e..5bc67633 100644 --- a/objectstore-service/src/backend/gcs.rs +++ b/objectstore-service/src/backend/gcs.rs @@ -518,8 +518,7 @@ impl GcsBackend { bucket, cogs, } = config; - // Incomplete GCS uploads do not contribute to storage COGS. - let change_stream = streams.build_with_upload_sessions(cogs.as_ref(), false); + let change_stream = streams.build(cogs.as_ref()); let token_provider = if endpoint.is_none() { Some(PrefetchingTokenProvider::gcp_auth(TOKEN_SCOPES).await?) @@ -3260,13 +3259,22 @@ mod tests { time_expires: Some(Timestamp::now() + Duration::from_secs(3600)), ..Default::default() }; + let before = Timestamp::now(); let token = backend .create_upload_session(&id, &metadata, nonzero(payload.len() as u64)) .await?; - assert!(producer.records().is_empty()); + let created = producer.records(); + assert_eq!(created.len(), 1); + assert_eq!(created[0].size, Some(payload.len() as u64)); + let expiration = created[0].expiration_time.unwrap() as u64; + assert!( + ((before + UPLOAD_SESSION_TTL).as_micros() + ..=(Timestamp::now() + UPLOAD_SESSION_TTL).as_micros()) + .contains(&expiration) + ); backend.upload_offset(&token).await?; - assert!(producer.records().is_empty()); + assert_eq!(producer.records().len(), 1); assert_eq!( backend @@ -3281,14 +3289,19 @@ mod tests { ); let records = producer.records(); - assert_eq!(records.len(), 1); - assert_eq!(records[0].op_type, OpType::Write); + assert_eq!(records.len(), 3); assert_eq!( - records[0].size, + records.iter().map(|r| r.op_type).collect::>(), + [OpType::Write, OpType::Write, OpType::Delete] + ); + assert_eq!(records[0].record_id, records[2].record_id); + assert_ne!(records[0].record_id, records[1].record_id); + assert_eq!( + records[1].size, Some(payload.len() as u64 + GcsObject::from_metadata(&metadata).metadata_size()) ); assert_eq!( - records[0].expiration_time, + records[1].expiration_time, metadata.time_expires.map(|t| t.as_micros() as i64) ); producer.clear(); @@ -3297,7 +3310,11 @@ mod tests { .await?; backend.cancel_upload(&canceled).await?; backend.cancel_upload(&canceled).await?; - assert!(producer.records().is_empty()); + let records = producer.records(); + assert_eq!(records.len(), 3); + assert_eq!(records[1].op_type, OpType::Delete); + assert_eq!(records[2].op_type, OpType::Delete); + assert!(records.iter().all(|r| r.record_id == records[0].record_id)); Ok(()) } @@ -3335,10 +3352,15 @@ mod tests { ); let records = producer.records(); - assert_eq!(records.len(), 1); - assert_eq!(records[0].op_type, OpType::Write); + assert_eq!(records.len(), 3); assert_eq!( - records[0].size, + records.iter().map(|r| r.op_type).collect::>(), + [OpType::Write, OpType::Write, OpType::Delete] + ); + assert_eq!(records[0].record_id, records[2].record_id); + assert_ne!(records[0].record_id, records[1].record_id); + assert_eq!( + records[1].size, Some(payload.len() as u64 + GcsObject::from_metadata(&metadata).metadata_size()) ); Ok(()) diff --git a/objectstore-service/src/backend/tiered.rs b/objectstore-service/src/backend/tiered.rs index 6d8dc19a..f458fcd8 100644 --- a/objectstore-service/src/backend/tiered.rs +++ b/objectstore-service/src/backend/tiered.rs @@ -1459,12 +1459,13 @@ mod tests { #[cfg(feature = "storage-cogs")] { let records = producer.records(); - assert_eq!(records.len(), 1); + assert_eq!(records.len(), 2); + assert_eq!(records[0].size, Some(payload.len() as u64)); assert_eq!( - records[0].size, + records[1].size, Some(revision.as_upload_path().to_string().len() as u64 + 1) ); - let expiry = records[0].expiration_time.unwrap() as u64; + let expiry = records[1].expiration_time.unwrap() as u64; assert!( ((before + RESUMABLE_UPLOAD_TTL).as_micros() ..=(Timestamp::now() + RESUMABLE_UPLOAD_TTL).as_micros()) @@ -1514,7 +1515,7 @@ mod tests { #[cfg(feature = "storage-cogs")] assert_eq!( producer.records().len(), - 1, + 2, "chunks and queries do not report changes" ); @@ -1556,14 +1557,18 @@ mod tests { .map(|r| (r.shared_resource_id.as_str(), r.op_type)) .collect::>(), [ + ("gcs_objectstore", Write), ("bigtable_objectstore", Write), ("bigtable_objectstore", Delete), ("gcs_objectstore", Write), + ("gcs_objectstore", Delete), ("bigtable_objectstore", Write), ] ); - assert_eq!(records[0].record_id, records[1].record_id); + assert_eq!(records[0].record_id, records[4].record_id); + assert_eq!(records[1].record_id, records[2].record_id); assert_ne!(records[0].record_id, records[3].record_id); + assert_ne!(records[1].record_id, records[5].record_id); producer.clear(); let canceled = @@ -1576,9 +1581,10 @@ mod tests { let records = producer.records(); assert_eq!( records.iter().map(|r| r.op_type).collect::>(), - [Write, Delete] + [Write, Write, Delete, Delete] ); - assert_eq!(records[0].record_id, records[1].record_id); + assert_eq!(records[0].record_id, records[3].record_id); + assert_eq!(records[1].record_id, records[2].record_id); } Ok(()) } diff --git a/objectstore-service/src/change_stream/cost_tracker.rs b/objectstore-service/src/change_stream/cost_tracker.rs index fc9fccc5..11bc7af5 100644 --- a/objectstore-service/src/change_stream/cost_tracker.rs +++ b/objectstore-service/src/change_stream/cost_tracker.rs @@ -25,14 +25,9 @@ use crate::id::ObjectId; /// can change the decision between a session write and delete. Object and session sampling /// decisions may differ. Neither raw paths nor session tokens are emitted. /// -/// Upload sessions are tracked by default. GCS disables their accounting at stream -/// construction because incomplete uploads do not incur storage charges. Generic -/// session lifecycle events remain available to other change stream implementations. -/// /// Logs, counts, and swallows errors returned by the [`InventoryTracker`]. pub struct CostTrackerStream { tracker: InventoryTracker

, - pub(super) track_upload_sessions: bool, } impl CostTrackerStream

{ @@ -44,7 +39,6 @@ impl CostTrackerStream

{ &config.shared_resource_id, config.sample_rate, ), - track_upload_sessions: true, } } @@ -78,7 +72,6 @@ impl fmt::Debug for CostTrackerStream

{ f.debug_struct("CostTrackerStream") .field("shared_resource_id", &self.tracker.shared_resource_id()) .field("sample_rate", &self.tracker.sample_rate()) - .field("track_upload_sessions", &self.track_upload_sessions) .finish() } } @@ -113,9 +106,6 @@ where P::Error: Into + Send + 'static, { fn write(&self, target: ChangeTarget<'_>, size: u64, expires_at: Option) { - if !self.track_upload_sessions && matches!(target, ChangeTarget::UploadSession { .. }) { - return; - } let (id, key) = target.inventory_identity(); let result = self.tracker.write( &key, @@ -130,9 +120,6 @@ where } fn update(&self, target: ChangeTarget<'_>, expires_at: Option) { - if !self.track_upload_sessions && matches!(target, ChangeTarget::UploadSession { .. }) { - return; - } let (id, key) = target.inventory_identity(); let result = self.tracker.update( &key, @@ -146,9 +133,6 @@ where } fn delete(&self, target: ChangeTarget<'_>) { - if !self.track_upload_sessions && matches!(target, ChangeTarget::UploadSession { .. }) { - return; - } let (id, key) = target.inventory_identity(); let result = self.tracker.delete(&key, id.usecase(), SystemTime::now()); self.swallow("delete", result); @@ -335,35 +319,6 @@ mod tests { assert!(!record_ids.contains(&producer.records()[0].record_id)); } - #[test] - fn upload_sessions_can_be_excluded_without_affecting_objects() { - let (producer, mut stream) = stream(1.0); - let id = object_id("attachments/org.17/project.42/objects/abc"); - let session = ChangeTarget::UploadSession { - object_id: &id, - session_id: "upload-token", - }; - stream.track_upload_sessions = false; - stream.write(session, 10, None); - stream.update(session, None); - stream.delete(session); - assert!(producer.records().is_empty()); - - stream.write((&id).into(), 10, None); - stream.update((&id).into(), None); - stream.delete((&id).into()); - let records = producer.records(); - assert_eq!(records.len(), 3); - assert_eq!(records[0].op_type, OpType::Write); - assert_eq!(records[1].op_type, OpType::Update); - assert_eq!(records[2].op_type, OpType::Delete); - assert!( - records - .iter() - .all(|record| record.record_id == records[0].record_id) - ); - } - #[test] fn a_listener_sampled_at_zero_reports_nothing() { let (producer, stream) = stream(0.0); diff --git a/objectstore-service/src/change_stream/factory.rs b/objectstore-service/src/change_stream/factory.rs index 2dced40c..5a69c610 100644 --- a/objectstore-service/src/change_stream/factory.rs +++ b/objectstore-service/src/change_stream/factory.rs @@ -48,23 +48,10 @@ impl ChangeStreamFactory { } /// Builds the stream `config` asks for, or a [`NoopStream`] if it cannot be built. - pub fn build(&self, config: Option<&CostTrackerStreamConfig>) -> Arc { - self.build_with_upload_sessions(config, true) - } - - /// Builds a stream with the backend's upload-session accounting policy. #[cfg(feature = "storage-cogs")] - pub(crate) fn build_with_upload_sessions( - &self, - config: Option<&CostTrackerStreamConfig>, - track_upload_sessions: bool, - ) -> Arc { + pub fn build(&self, config: Option<&CostTrackerStreamConfig>) -> Arc { match (config, self.producer.clone()) { - (Some(config), Some(producer)) => { - let mut stream = CostTrackerStream::new(producer, config); - stream.track_upload_sessions = track_upload_sessions; - Arc::new(stream) - } + (Some(config), Some(producer)) => Arc::new(CostTrackerStream::new(producer, config)), (None, None) => Arc::new(NoopStream), (c, p) => { objectstore_log::warn!( @@ -79,11 +66,7 @@ impl ChangeStreamFactory { /// Reporting is not compiled in, so every backend reports nothing. #[cfg(not(feature = "storage-cogs"))] - pub(crate) fn build_with_upload_sessions( - &self, - _config: Option<&CostTrackerStreamConfig>, - _track_upload_sessions: bool, - ) -> Arc { + pub fn build(&self, _config: Option<&CostTrackerStreamConfig>) -> Arc { Arc::new(NoopStream) } } diff --git a/objectstore-service/src/change_stream/mod.rs b/objectstore-service/src/change_stream/mod.rs index 3a88d629..9747f427 100644 --- a/objectstore-service/src/change_stream/mod.rs +++ b/objectstore-service/src/change_stream/mod.rs @@ -4,9 +4,6 @@ //! the service describes where those records go with a [`CostTrackerConfig`], shared by //! every backend. [`ChangeStreamFactory`] pairs the two into a [`ChangeStream`]. //! -//! Cost tracking includes upload sessions by default. GCS disables session accounting in -//! the cost-tracking adapter while still publishing the generic lifecycle events. -//! //! # Resumable uploads //! //! Upload backends report a session write after creation with exactly the advertised upload From 250bed0b07d07aee0781c6f22c1ac510f55cf13c Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Thu, 8 Oct 2026 16:13:43 +0200 Subject: [PATCH 04/11] feat(cogs): Distinguish upload session operations on the wire Emit WRITE_SESSION and DELETE_SESSION for upload estimates while keeping stored upload markers on WRITE and DELETE. Preserve existing record identities and add session-specific inventory tracker methods without changing callers. --- objectstore-inventory-tracker/README.md | 6 ++ objectstore-inventory-tracker/src/lib.rs | 8 +- objectstore-inventory-tracker/src/record.rs | 19 ++-- objectstore-inventory-tracker/src/tracker.rs | 93 ++++++++++++++++++- objectstore-service/docs/architecture.md | 12 ++- objectstore-service/src/backend/gcs.rs | 8 +- objectstore-service/src/backend/in_memory.rs | 6 +- objectstore-service/src/backend/local_fs.rs | 6 +- objectstore-service/src/backend/tiered.rs | 10 +- .../src/change_stream/cost_tracker.rs | 48 ++++++---- objectstore-service/src/change_stream/mod.rs | 16 ++-- 11 files changed, 180 insertions(+), 52 deletions(-) diff --git a/objectstore-inventory-tracker/README.md b/objectstore-inventory-tracker/README.md index 476d0138..2c207b47 100644 --- a/objectstore-inventory-tracker/README.md +++ b/objectstore-inventory-tracker/README.md @@ -56,3 +56,9 @@ recommended that you configure each of them with their own `InventoryTracker`. |---|---| | `kafka` | The `sentry_arroyo`-backed producer. Off by default, so the record types are usable without pulling in arroyo, librdkafka, and their native build. | | `test-utils` | Exposes the `test_utils` module and its `DummyProducer`, so downstream crates can assert on emitted records without a broker. | + +Upload sessions can be reported with `InventoryTracker::write_session` and +`InventoryTracker::delete_session`. They emit `WRITE_SESSION` and `DELETE_SESSION` +so downstream consumers can exclude session estimates from storage accounting. +Use a stable session key distinct from the final object's key. Deploy the updated +Kafka schema to validating consumers before emitting these operation types. diff --git a/objectstore-inventory-tracker/src/lib.rs b/objectstore-inventory-tracker/src/lib.rs index f7ed4dbe..e35f7f76 100644 --- a/objectstore-inventory-tracker/src/lib.rs +++ b/objectstore-inventory-tracker/src/lib.rs @@ -6,7 +6,9 @@ //! is meant to be, for example, a specific GCS bucket or Bigtable instance. //! - `record_id`: identifies each record. `InventoryTracker` populates this with a hash //! of the identifier passed in by the caller. -//! - `size`: the size of the record in bytes (including metadata). +//! - `op_type`: `WRITE`, `UPDATE`, or `DELETE` for stored objects; `WRITE_SESSION` or +//! `DELETE_SESSION` for upload sessions that consumers may exclude from accounting. +//! - `size`: stored bytes (including metadata), or the estimated upload-session size. //! - `expiration_time`: a timestamp (unixtime microseconds) describing when the record is //! meant to be deleted. //! @@ -34,6 +36,10 @@ //! analyze costs. If you have multiple buckets, or multiple storage backends, it's //! recommended that you configure each of them with their own `InventoryTracker`. //! +//! Use [`InventoryTracker::write_session`] and [`InventoryTracker::delete_session`] for +//! upload sessions, with a stable key distinct from the final object's key. Consumers +//! must support the session operation types before producers start emitting them. +//! //! # Sampling //! //! [`InventoryTracker`] hashes the storage key it gets from the caller and uses part of diff --git a/objectstore-inventory-tracker/src/record.rs b/objectstore-inventory-tracker/src/record.rs index bb3378c7..b535e80d 100644 --- a/objectstore-inventory-tracker/src/record.rs +++ b/objectstore-inventory-tracker/src/record.rs @@ -1,9 +1,8 @@ //! The wire format emitted onto the inventory topic. //! //! These types mirror the `shared-resources-inventory` schema registered in -//! [sentry-kafka-schemas]. The schema sets `additionalProperties: false`, so adding a -//! field here without a corresponding schema version bump produces messages that -//! consumers reject. +//! [sentry-kafka-schemas]. New fields or operation types require a coordinated schema +//! update and rollout to validating consumers before producers emit them. //! //! [sentry-kafka-schemas]: https://github.com/getsentry/sentry-kafka-schemas @@ -11,9 +10,9 @@ use std::time::{SystemTime, UNIX_EPOCH}; use serde::{Deserialize, Serialize}; -/// The kind of change a record describes. `WRITE`, `UPDATE`, or `DELETE`. +/// The kind of object or upload-session change a record describes. #[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)] -#[serde(rename_all = "UPPERCASE")] +#[serde(rename_all = "SCREAMING_SNAKE_CASE")] pub enum OpType { /// The record was created, or replaced with new contents. Write, @@ -23,6 +22,10 @@ pub enum OpType { Update, /// The record was removed. Delete, + /// An upload session was created, or its estimated size was replaced. + WriteSession, + /// An upload session completed or was canceled. + DeleteSession, } /// A single inventory change event. @@ -61,7 +64,9 @@ pub struct InventoryRecord { /// sample rate. pub sample_rate: f64, - /// Stored size in bytes. Always set for [`OpType::Write`]. + /// Stored or estimated size in bytes. + /// + /// Always set for [`OpType::Write`] and [`OpType::WriteSession`]. #[serde(skip_serializing_if = "Option::is_none")] pub size: Option, @@ -162,6 +167,8 @@ mod tests { (OpType::Write, "WRITE"), (OpType::Update, "UPDATE"), (OpType::Delete, "DELETE"), + (OpType::WriteSession, "WRITE_SESSION"), + (OpType::DeleteSession, "DELETE_SESSION"), ] { assert_eq!(serde_json::to_value(op_type).unwrap(), json!(expected)); } diff --git a/objectstore-inventory-tracker/src/tracker.rs b/objectstore-inventory-tracker/src/tracker.rs index 7897012b..08a3cd2a 100644 --- a/objectstore-inventory-tracker/src/tracker.rs +++ b/objectstore-inventory-tracker/src/tracker.rs @@ -153,6 +153,58 @@ impl InventoryTracker

{ expiration_time: Option, organization_id: Option, project_id: Option, + ) -> Result<(), P::Error> { + self.write_with_op( + OpType::Write, + storage_key, + app_feature, + size, + timestamp, + expiration_time, + organization_id, + project_id, + ) + } + + /// Emits a `WRITE_SESSION`: an upload session now occupies the estimated size. + /// + /// Use a stable session key distinct from the eventual object's key. + /// + /// Does nothing and returns `Ok(())` if `storage_key` is not sampled. + #[allow(clippy::too_many_arguments)] + pub fn write_session( + &self, + storage_key: &str, + app_feature: &str, + size: u64, + timestamp: SystemTime, + expiration_time: Option, + organization_id: Option, + project_id: Option, + ) -> Result<(), P::Error> { + self.write_with_op( + OpType::WriteSession, + storage_key, + app_feature, + size, + timestamp, + expiration_time, + organization_id, + project_id, + ) + } + + #[allow(clippy::too_many_arguments)] + fn write_with_op( + &self, + op_type: OpType, + storage_key: &str, + app_feature: &str, + size: u64, + timestamp: SystemTime, + expiration_time: Option, + organization_id: Option, + project_id: Option, ) -> Result<(), P::Error> { let Some(record_id) = self.sample(storage_key) else { return Ok(()); @@ -161,7 +213,7 @@ impl InventoryTracker

{ self.emit(InventoryRecord { shared_resource_id: self.shared_resource_id.clone(), app_feature: app_feature.to_owned(), - op_type: OpType::Write, + op_type, record_id, timestamp: epoch_micros(timestamp), sample_rate: self.sample_rate, @@ -211,6 +263,28 @@ impl InventoryTracker

{ storage_key: &str, app_feature: &str, timestamp: SystemTime, + ) -> Result<(), P::Error> { + self.delete_with_op(OpType::Delete, storage_key, app_feature, timestamp) + } + + /// Emits a `DELETE_SESSION`: an upload session completed or was canceled. + /// + /// Does nothing and returns `Ok(())` if `storage_key` is not sampled. + pub fn delete_session( + &self, + storage_key: &str, + app_feature: &str, + timestamp: SystemTime, + ) -> Result<(), P::Error> { + self.delete_with_op(OpType::DeleteSession, storage_key, app_feature, timestamp) + } + + fn delete_with_op( + &self, + op_type: OpType, + storage_key: &str, + app_feature: &str, + timestamp: SystemTime, ) -> Result<(), P::Error> { let Some(record_id) = self.sample(storage_key) else { return Ok(()); @@ -219,7 +293,7 @@ impl InventoryTracker

{ self.emit(InventoryRecord { shared_resource_id: self.shared_resource_id.clone(), app_feature: app_feature.to_owned(), - op_type: OpType::Delete, + op_type, record_id, timestamp: epoch_micros(timestamp), sample_rate: self.sample_rate, @@ -393,6 +467,10 @@ mod tests { .update("some/key", "f", now, Some(now), None, None) .unwrap(); tracker.delete("some/key", "f", now).unwrap(); + tracker + .write_session("session/key", "f", 4096, now, Some(now), Some(1), Some(2)) + .unwrap(); + tracker.delete_session("session/key", "f", now).unwrap(); let records = producer.records(); assert_eq!(records[0].op_type, OpType::Write); @@ -401,6 +479,17 @@ mod tests { assert_eq!(records[1].size, None, "update means size unchanged"); assert_eq!(records[2].op_type, OpType::Delete); assert_eq!(records[2].size, None); + assert_eq!(records.len(), 5); + assert_eq!(records[3].op_type, OpType::WriteSession); + assert_eq!(records[3].size, Some(4096)); + assert_eq!(records[3].expiration_time, Some(epoch_micros(now))); + assert_eq!(records[3].organization_id, Some(1)); + assert_eq!(records[3].project_id, Some(2)); + assert_eq!(records[4].op_type, OpType::DeleteSession); + assert_eq!(records[4].size, None); + assert_eq!(records[4].expiration_time, None); + assert_eq!(records[3].record_id, records[4].record_id); + assert_ne!(records[0].record_id, records[3].record_id); } #[test] diff --git a/objectstore-service/docs/architecture.md b/objectstore-service/docs/architecture.md index 2d7e5d9d..d432062b 100644 --- a/objectstore-service/docs/architecture.md +++ b/objectstore-service/docs/architecture.md @@ -164,9 +164,11 @@ backend. When using [`TieredStorage`](backend::tiered::TieredStorage)'s long-term backend the inventory table will contain _two rows_ for an object: a row for the actual object and its size in long-term backend, and a separate row for the tombstone and the tombstone's size in the high-volume backend. -While a resumable upload is in progress, separate session rows account for its -advertised size in the upload backend and its marker's stored size in high-volume -storage. See the [change stream module](change_stream) for their lifetimes. +While a resumable upload is in progress, `WRITE_SESSION` and `DELETE_SESSION` events +track its advertised size separately from the published object. Consumers can exclude +these estimates from cost attribution. High-volume markers still use ordinary +`WRITE` and `DELETE` events for their actual stored bytes. +See the [change stream module](change_stream) for their lifetimes. Because the change stream does not observe automatic garbage collection, expired records must be filtered out when querying the inventory table. @@ -197,8 +199,8 @@ per-backend feed of three operations: unchanged. In practice this is a TTI bump. - `delete(target)`: `target` was deleted explicitly. -[`ChangeTarget`](change_stream::ChangeTarget) identifies either an object or an upload -session. An upload session carries its object's identity for attribution and a stable +[`ChangeTarget`](change_stream::ChangeTarget) identifies an object, a stored upload +marker, or an upload session. An upload session carries its object's identity for attribution and a stable session ID; [`Session`](resumable::Session) converts directly into a session target. The stream describes physical storage per backend. When using diff --git a/objectstore-service/src/backend/gcs.rs b/objectstore-service/src/backend/gcs.rs index 5bc67633..06f1956a 100644 --- a/objectstore-service/src/backend/gcs.rs +++ b/objectstore-service/src/backend/gcs.rs @@ -3292,7 +3292,7 @@ mod tests { assert_eq!(records.len(), 3); assert_eq!( records.iter().map(|r| r.op_type).collect::>(), - [OpType::Write, OpType::Write, OpType::Delete] + [OpType::WriteSession, OpType::Write, OpType::DeleteSession] ); assert_eq!(records[0].record_id, records[2].record_id); assert_ne!(records[0].record_id, records[1].record_id); @@ -3312,8 +3312,8 @@ mod tests { backend.cancel_upload(&canceled).await?; let records = producer.records(); assert_eq!(records.len(), 3); - assert_eq!(records[1].op_type, OpType::Delete); - assert_eq!(records[2].op_type, OpType::Delete); + assert_eq!(records[1].op_type, OpType::DeleteSession); + assert_eq!(records[2].op_type, OpType::DeleteSession); assert!(records.iter().all(|r| r.record_id == records[0].record_id)); Ok(()) } @@ -3355,7 +3355,7 @@ mod tests { assert_eq!(records.len(), 3); assert_eq!( records.iter().map(|r| r.op_type).collect::>(), - [OpType::Write, OpType::Write, OpType::Delete] + [OpType::WriteSession, OpType::Write, OpType::DeleteSession] ); assert_eq!(records[0].record_id, records[2].record_id); assert_ne!(records[0].record_id, records[1].record_id); diff --git a/objectstore-service/src/backend/in_memory.rs b/objectstore-service/src/backend/in_memory.rs index 6650f349..c4698a8c 100644 --- a/objectstore-service/src/backend/in_memory.rs +++ b/objectstore-service/src/backend/in_memory.rs @@ -1168,7 +1168,7 @@ mod tests { assert_eq!(records.len(), 3); assert_eq!( records.iter().map(|r| r.op_type).collect::>(), - [OpType::Write, OpType::Write, OpType::Delete] + [OpType::WriteSession, OpType::Write, OpType::DeleteSession] ); assert_eq!(records[0].record_id, records[2].record_id); assert_ne!(records[0].record_id, records[1].record_id); @@ -1192,8 +1192,8 @@ mod tests { assert_ne!(records[0].record_id, records[1].record_id); assert_eq!(records[0].record_id, records[2].record_id); assert_eq!(records[1].record_id, records[3].record_id); - assert_eq!(records[2].op_type, OpType::Delete); - assert_eq!(records[3].op_type, OpType::Delete); + assert_eq!(records[2].op_type, OpType::DeleteSession); + assert_eq!(records[3].op_type, OpType::DeleteSession); } #[tokio::test] diff --git a/objectstore-service/src/backend/local_fs.rs b/objectstore-service/src/backend/local_fs.rs index 4c6f0935..9c4d4fc9 100644 --- a/objectstore-service/src/backend/local_fs.rs +++ b/objectstore-service/src/backend/local_fs.rs @@ -1280,7 +1280,7 @@ mod tests { let records = producer.records(); assert_eq!( records.iter().map(|r| r.op_type).collect::>(), - [OpType::Write, OpType::Write, OpType::Delete] + [OpType::WriteSession, OpType::Write, OpType::DeleteSession] ); assert_eq!(records[0].record_id, records[2].record_id); } @@ -2345,7 +2345,7 @@ mod tests { assert_eq!(records.len(), 3); assert_eq!( records.iter().map(|r| r.op_type).collect::>(), - [OpType::Write, OpType::Write, OpType::Delete] + [OpType::WriteSession, OpType::Write, OpType::DeleteSession] ); assert_eq!(records[0].record_id, records[2].record_id); assert_ne!(records[0].record_id, records[1].record_id); @@ -2361,7 +2361,7 @@ mod tests { assert!(backend.cancel_upload(&canceled).await.is_err()); let records = producer.records(); assert_eq!(records.len(), 2); - assert_eq!(records[1].op_type, OpType::Delete); + assert_eq!(records[1].op_type, OpType::DeleteSession); assert_eq!(records[0].record_id, records[1].record_id); } diff --git a/objectstore-service/src/backend/tiered.rs b/objectstore-service/src/backend/tiered.rs index f458fcd8..02fd2a28 100644 --- a/objectstore-service/src/backend/tiered.rs +++ b/objectstore-service/src/backend/tiered.rs @@ -1549,7 +1549,9 @@ mod tests { assert_eq!(stream::read_to_vec(body).await?, payload); #[cfg(feature = "storage-cogs")] { - use objectstore_inventory_tracker::OpType::{Delete, Write}; + use objectstore_inventory_tracker::OpType::{ + Delete, DeleteSession, Write, WriteSession, + }; let records = producer.records(); assert_eq!( records @@ -1557,11 +1559,11 @@ mod tests { .map(|r| (r.shared_resource_id.as_str(), r.op_type)) .collect::>(), [ - ("gcs_objectstore", Write), + ("gcs_objectstore", WriteSession), ("bigtable_objectstore", Write), ("bigtable_objectstore", Delete), ("gcs_objectstore", Write), - ("gcs_objectstore", Delete), + ("gcs_objectstore", DeleteSession), ("bigtable_objectstore", Write), ] ); @@ -1581,7 +1583,7 @@ mod tests { let records = producer.records(); assert_eq!( records.iter().map(|r| r.op_type).collect::>(), - [Write, Write, Delete, Delete] + [WriteSession, Write, Delete, DeleteSession] ); assert_eq!(records[0].record_id, records[3].record_id); assert_eq!(records[1].record_id, records[2].record_id); diff --git a/objectstore-service/src/change_stream/cost_tracker.rs b/objectstore-service/src/change_stream/cost_tracker.rs index 11bc7af5..bfcd2718 100644 --- a/objectstore-service/src/change_stream/cost_tracker.rs +++ b/objectstore-service/src/change_stream/cost_tracker.rs @@ -25,6 +25,10 @@ use crate::id::ObjectId; /// can change the decision between a session write and delete. Object and session sampling /// decisions may differ. Neither raw paths nor session tokens are emitted. /// +/// Sessions emit `WRITE_SESSION` and `DELETE_SESSION`. High-volume upload markers retain +/// their existing identities but emit ordinary `WRITE` and `DELETE`, since their stored +/// bytes contribute to cost attribution. +/// /// Logs, counts, and swallows errors returned by the [`InventoryTracker`]. pub struct CostTrackerStream { tracker: InventoryTracker

, @@ -79,23 +83,20 @@ impl fmt::Debug for CostTrackerStream

{ impl<'a> ChangeTarget<'a> { /// Returns the attribution ID and stable input to inventory hashing and sampling. fn inventory_identity(self) -> (&'a ObjectId, String) { - match self { - Self::Object(id) => (id, id.as_storage_path().to_string()), + let (object_id, session_id) = match self { + Self::Object(id) => return (id, id.as_storage_path().to_string()), + Self::UploadMarker(revision) => (revision, revision.key.as_str()), Self::UploadSession { object_id, session_id, - } => { - // This tuple is the permanent session identity contract, just as the - // storage path is for objects. Neither the path nor token is emitted. - let key = serde_json::to_string(&( - "upload_session", - object_id.as_storage_path(), - session_id, - )) + } => (object_id, session_id), + }; + // Preserve the established identity for both sessions and markers. The operation + // type distinguishes session estimates from stored markers on the wire. + let key = + serde_json::to_string(&("upload_session", object_id.as_storage_path(), session_id)) .expect("session identity is serializable"); - (object_id, key) - } - } + (object_id, key) } } @@ -107,7 +108,12 @@ where { fn write(&self, target: ChangeTarget<'_>, size: u64, expires_at: Option) { let (id, key) = target.inventory_identity(); - let result = self.tracker.write( + let write = match target { + ChangeTarget::Object(_) | ChangeTarget::UploadMarker(_) => InventoryTracker::write, + ChangeTarget::UploadSession { .. } => InventoryTracker::write_session, + }; + let result = write( + &self.tracker, &key, id.usecase(), size, @@ -134,7 +140,15 @@ where fn delete(&self, target: ChangeTarget<'_>) { let (id, key) = target.inventory_identity(); - let result = self.tracker.delete(&key, id.usecase(), SystemTime::now()); + let result = match target { + ChangeTarget::Object(_) | ChangeTarget::UploadMarker(_) => { + self.tracker.delete(&key, id.usecase(), SystemTime::now()) + } + ChangeTarget::UploadSession { .. } => { + self.tracker + .delete_session(&key, id.usecase(), SystemTime::now()) + } + }; self.swallow("delete", result); } @@ -301,8 +315,8 @@ mod tests { if !records.is_empty() { sampled += 1; assert_eq!(records.len(), 2); - assert_eq!(records[0].op_type, OpType::Write); - assert_eq!(records[1].op_type, OpType::Delete); + assert_eq!(records[0].op_type, OpType::WriteSession); + assert_eq!(records[1].op_type, OpType::DeleteSession); assert_eq!(records[0].record_id, records[1].record_id); assert!(record_ids.insert(records[0].record_id.clone())); assert_eq!(records[0].organization_id, Some(17)); diff --git a/objectstore-service/src/change_stream/mod.rs b/objectstore-service/src/change_stream/mod.rs index 9747f427..cc3b3a04 100644 --- a/objectstore-service/src/change_stream/mod.rs +++ b/objectstore-service/src/change_stream/mod.rs @@ -11,9 +11,12 @@ //! expiration. Partial chunks and incomplete offset queries emit nothing. Successful publication //! reports the object's actual stored size and expiration, followed by a session delete. Completion //! discovered through an offset query follows the same order. Successful cancellation also -//! reports a session delete. The accounting deadline does not change backend cleanup behavior. +//! reports a session delete. Session writes and deletes use `WRITE_SESSION` and `DELETE_SESSION` +//! on the inventory wire, allowing consumers to exclude their estimates from cost attribution. +//! The accounting deadline does not change backend cleanup behavior. //! -//! High-volume markers report their stored size and actual deadline, set by +//! High-volume markers use ordinary `WRITE`/`DELETE` operations to report their stored size +//! and actual deadline, set by //! `backend::tiered::RESUMABLE_UPLOAD_TTL`. Only a successful //! conditional marker deletion emits a delete. Tiered deletes its marker to claim the upload //! before finalization, so that marker delete precedes the long-term object write and session @@ -61,7 +64,9 @@ pub(crate) const UPLOAD_SESSION_TTL: Duration = Duration::from_hours(7 * 24); pub enum ChangeTarget<'a> { /// A published object, including a high-volume redirect tombstone. Object(&'a ObjectId), - /// An in-progress upload or its high-volume marker. + /// A stored high-volume upload marker, accounted for as an ordinary object. + UploadMarker(&'a ObjectId), + /// An in-progress upload, identified by its backend session token. UploadSession { /// Object identity supplying the usecase and scopes. object_id: &'a ObjectId, @@ -73,10 +78,7 @@ pub enum ChangeTarget<'a> { impl<'a> ChangeTarget<'a> { /// Identifies the high-volume marker for an upload's unique revision. pub fn upload_marker(revision: &'a ObjectId) -> Self { - Self::UploadSession { - object_id: revision, - session_id: &revision.key, - } + Self::UploadMarker(revision) } } From be544a0d86cfd3c047e3ff1234846a6b1f317d92 Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Thu, 8 Oct 2026 16:48:56 +0200 Subject: [PATCH 05/11] improve --- objectstore-inventory-tracker/README.md | 6 ------ objectstore-inventory-tracker/src/lib.rs | 6 +----- objectstore-inventory-tracker/src/record.rs | 5 +++-- objectstore-service/docs/architecture.md | 13 +++---------- objectstore-service/src/backend/gcs.rs | 14 +++++++++++--- objectstore-service/src/backend/in_memory.rs | 10 ++++++++-- objectstore-service/src/backend/local_fs.rs | 17 +++++++++++------ .../src/change_stream/cost_tracker.rs | 14 +------------- objectstore-service/src/change_stream/mod.rs | 12 ++++-------- 9 files changed, 42 insertions(+), 55 deletions(-) diff --git a/objectstore-inventory-tracker/README.md b/objectstore-inventory-tracker/README.md index 2c207b47..476d0138 100644 --- a/objectstore-inventory-tracker/README.md +++ b/objectstore-inventory-tracker/README.md @@ -56,9 +56,3 @@ recommended that you configure each of them with their own `InventoryTracker`. |---|---| | `kafka` | The `sentry_arroyo`-backed producer. Off by default, so the record types are usable without pulling in arroyo, librdkafka, and their native build. | | `test-utils` | Exposes the `test_utils` module and its `DummyProducer`, so downstream crates can assert on emitted records without a broker. | - -Upload sessions can be reported with `InventoryTracker::write_session` and -`InventoryTracker::delete_session`. They emit `WRITE_SESSION` and `DELETE_SESSION` -so downstream consumers can exclude session estimates from storage accounting. -Use a stable session key distinct from the final object's key. Deploy the updated -Kafka schema to validating consumers before emitting these operation types. diff --git a/objectstore-inventory-tracker/src/lib.rs b/objectstore-inventory-tracker/src/lib.rs index e35f7f76..9ab79877 100644 --- a/objectstore-inventory-tracker/src/lib.rs +++ b/objectstore-inventory-tracker/src/lib.rs @@ -7,7 +7,7 @@ //! - `record_id`: identifies each record. `InventoryTracker` populates this with a hash //! of the identifier passed in by the caller. //! - `op_type`: `WRITE`, `UPDATE`, or `DELETE` for stored objects; `WRITE_SESSION` or -//! `DELETE_SESSION` for upload sessions that consumers may exclude from accounting. +//! `DELETE_SESSION` for upload sessions. //! - `size`: stored bytes (including metadata), or the estimated upload-session size. //! - `expiration_time`: a timestamp (unixtime microseconds) describing when the record is //! meant to be deleted. @@ -36,10 +36,6 @@ //! analyze costs. If you have multiple buckets, or multiple storage backends, it's //! recommended that you configure each of them with their own `InventoryTracker`. //! -//! Use [`InventoryTracker::write_session`] and [`InventoryTracker::delete_session`] for -//! upload sessions, with a stable key distinct from the final object's key. Consumers -//! must support the session operation types before producers start emitting them. -//! //! # Sampling //! //! [`InventoryTracker`] hashes the storage key it gets from the caller and uses part of diff --git a/objectstore-inventory-tracker/src/record.rs b/objectstore-inventory-tracker/src/record.rs index b535e80d..aacda9d8 100644 --- a/objectstore-inventory-tracker/src/record.rs +++ b/objectstore-inventory-tracker/src/record.rs @@ -1,8 +1,9 @@ //! The wire format emitted onto the inventory topic. //! //! These types mirror the `shared-resources-inventory` schema registered in -//! [sentry-kafka-schemas]. New fields or operation types require a coordinated schema -//! update and rollout to validating consumers before producers emit them. +//! [sentry-kafka-schemas]. The schema sets `additionalProperties: false`, so adding a +//! field here without a corresponding schema version bump produces messages that +//! consumers reject. //! //! [sentry-kafka-schemas]: https://github.com/getsentry/sentry-kafka-schemas diff --git a/objectstore-service/docs/architecture.md b/objectstore-service/docs/architecture.md index d432062b..7cb2ef1b 100644 --- a/objectstore-service/docs/architecture.md +++ b/objectstore-service/docs/architecture.md @@ -164,11 +164,6 @@ backend. When using [`TieredStorage`](backend::tiered::TieredStorage)'s long-term backend the inventory table will contain _two rows_ for an object: a row for the actual object and its size in long-term backend, and a separate row for the tombstone and the tombstone's size in the high-volume backend. -While a resumable upload is in progress, `WRITE_SESSION` and `DELETE_SESSION` events -track its advertised size separately from the published object. Consumers can exclude -these estimates from cost attribution. High-volume markers still use ordinary -`WRITE` and `DELETE` events for their actual stored bytes. -See the [change stream module](change_stream) for their lifetimes. Because the change stream does not observe automatic garbage collection, expired records must be filtered out when querying the inventory table. @@ -200,8 +195,7 @@ per-backend feed of three operations: - `delete(target)`: `target` was deleted explicitly. [`ChangeTarget`](change_stream::ChangeTarget) identifies an object, a stored upload -marker, or an upload session. An upload session carries its object's identity for attribution and a stable -session ID; [`Session`](resumable::Session) converts directly into a session target. +marker, or an upload session. The stream describes physical storage per backend. When using [`TieredStorage`](backend::tiered::TieredStorage), objects that are stored in @@ -210,9 +204,8 @@ storage as well as for the tombstone record in high-volume storage. For objects and markers, `size` is a count of bytes that the backend actually stores. This includes object payloads, metadata, and sometimes backend-specific overhead. - -See the [change stream module](change_stream) for upload-session accounting, expiration, -and lifecycle reporting. +For upload sessions, `size` is the final `size` of the corresponding `object` that +the upload will create when the upload is finalized. Decorators such as [`CountingBackend`](backend::counting::CountingBackend) and [`TieredStorage`](backend::tiered::TieredStorage) don't publish change streams diff --git a/objectstore-service/src/backend/gcs.rs b/objectstore-service/src/backend/gcs.rs index 06f1956a..48cbd294 100644 --- a/objectstore-service/src/backend/gcs.rs +++ b/objectstore-service/src/backend/gcs.rs @@ -1156,7 +1156,9 @@ impl Backend for GcsBackend { ) -> Result> { objectstore_log::debug!("Creating resumable upload session on GCS backend"); let url = self.upload_url(id, "resumable")?; - let metadata_json = serde_json::to_vec(&GcsObject::from_metadata(metadata)).context( + let gcs_metadata = GcsObject::from_metadata(metadata); + let metadata_size = gcs_metadata.metadata_size(); + let metadata_json = serde_json::to_vec(&gcs_metadata).context( ErrorKind::Internal, "serializing GCS resumable upload metadata", )?; @@ -1218,7 +1220,7 @@ impl Backend for GcsBackend { object_id: id, session_id: &token, }, - upload_length.get(), + upload_length.get().saturating_add(metadata_size), Some(Timestamp::now() + UPLOAD_SESSION_TTL), ); Ok(Some(token)) @@ -3257,6 +3259,8 @@ mod tests { let payload = b"resumable payload".to_vec(); let metadata = Metadata { time_expires: Some(Timestamp::now() + Duration::from_secs(3600)), + filename: Some("upload.txt".into()), + custom: BTreeMap::from_iter([("hello".into(), "world".into())]), ..Default::default() }; let before = Timestamp::now(); @@ -3266,7 +3270,10 @@ mod tests { let created = producer.records(); assert_eq!(created.len(), 1); - assert_eq!(created[0].size, Some(payload.len() as u64)); + assert_eq!( + created[0].size, + Some(payload.len() as u64 + GcsObject::from_metadata(&metadata).metadata_size()) + ); let expiration = created[0].expiration_time.unwrap() as u64; assert!( ((before + UPLOAD_SESSION_TTL).as_micros() @@ -3296,6 +3303,7 @@ mod tests { ); assert_eq!(records[0].record_id, records[2].record_id); assert_ne!(records[0].record_id, records[1].record_id); + assert_eq!(records[0].size, records[1].size); assert_eq!( records[1].size, Some(payload.len() as u64 + GcsObject::from_metadata(&metadata).metadata_size()) diff --git a/objectstore-service/src/backend/in_memory.rs b/objectstore-service/src/backend/in_memory.rs index c4698a8c..01aa53b9 100644 --- a/objectstore-service/src/backend/in_memory.rs +++ b/objectstore-service/src/backend/in_memory.rs @@ -277,7 +277,9 @@ impl super::common::Backend for InMemoryBackend { object_id: id, session_id: &token, }, - upload_length.get(), + upload_length + .get() + .saturating_add(json_len(metadata) as u64), Some(Timestamp::now() + UPLOAD_SESSION_TTL), ); Ok(Some(token)) @@ -1144,7 +1146,10 @@ mod tests { let token = create_session(&backend, &id, 2).await; let created = producer.records(); assert_eq!(created.len(), 1); - assert_eq!(created[0].size, Some(2)); + assert_eq!( + created[0].size, + Some(2 + json_len(&Metadata::default()) as u64) + ); let expiration = created[0].expiration_time.unwrap() as u64; assert!( ((before + UPLOAD_SESSION_TTL).as_micros() @@ -1172,6 +1177,7 @@ mod tests { ); assert_eq!(records[0].record_id, records[2].record_id); assert_ne!(records[0].record_id, records[1].record_id); + assert_eq!(records[0].size, records[1].size); assert_eq!( records[1].size, Some((json_len(&Metadata::default()) + 2) as u64) diff --git a/objectstore-service/src/backend/local_fs.rs b/objectstore-service/src/backend/local_fs.rs index 9c4d4fc9..2fe52565 100644 --- a/objectstore-service/src/backend/local_fs.rs +++ b/objectstore-service/src/backend/local_fs.rs @@ -325,14 +325,14 @@ impl Backend for LocalFsBackend { ErrorKind::BackendFailure, "creating local-fs object directory", )?; - UploadFile::create(&path, metadata).await?; + let metadata_size = UploadFile::create(&path, metadata).await?; let token = upload_id.to_string(); self.change_stream.write( ChangeTarget::UploadSession { object_id: id, session_id: &token, }, - upload_length.get(), + upload_length.get().saturating_add(metadata_size), Some(Timestamp::now() + UPLOAD_SESSION_TTL), ); Ok(Some(token)) @@ -835,7 +835,7 @@ struct UploadFile { } impl UploadFile { - async fn create(path: &Path, metadata: &Metadata) -> Result<()> { + async fn create(path: &Path, metadata: &Metadata) -> Result { let mut options = OpenOptions::from(file_options()); options.create_new(true).read(true).write(true); @@ -843,11 +843,12 @@ impl UploadFile { ErrorKind::BackendFailure, "creating local-fs resumable upload", )?; - write_metadata_preamble(&mut file, metadata).await?; + let metadata_size = write_metadata_preamble(&mut file, metadata).await?; file.sync_data().await.context( ErrorKind::BackendFailure, "syncing local-fs resumable upload", - ) + )?; + Ok(metadata_size) } async fn open(path: &Path) -> Result { @@ -2310,7 +2311,10 @@ mod tests { let created = producer.records(); assert_eq!(created.len(), 1); - assert_eq!(created[0].size, Some(payload.len() as u64)); + let upload_path = backend.upload_path(Uuid::parse_str(&token.backend_token).unwrap()); + let metadata_size = tokio::fs::metadata(upload_path).await.unwrap().len(); + assert!(metadata_size > 0); + assert_eq!(created[0].size, Some(payload.len() as u64 + metadata_size)); let expiration = created[0].expiration_time.unwrap() as u64; assert!( ((before + UPLOAD_SESSION_TTL).as_micros() @@ -2350,6 +2354,7 @@ mod tests { assert_eq!(records[0].record_id, records[2].record_id); assert_ne!(records[0].record_id, records[1].record_id); assert_eq!(records[1].size, Some(file.len() as u64)); + assert_eq!(records[0].size, records[1].size); assert_eq!( records[1].expiration_time, metadata.time_expires.map(|t| t.as_micros() as i64) diff --git a/objectstore-service/src/change_stream/cost_tracker.rs b/objectstore-service/src/change_stream/cost_tracker.rs index bfcd2718..d45f27d5 100644 --- a/objectstore-service/src/change_stream/cost_tracker.rs +++ b/objectstore-service/src/change_stream/cost_tracker.rs @@ -14,21 +14,9 @@ use objectstore_types::time::Timestamp; use crate::id::ObjectId; /// Reports through an [`InventoryTracker`], which hashes each target identity both to -/// anonymize it and to decide whether it is sampled. Session identities include the object -/// path and session ID, keeping their records and sampling separate from published objects. +/// anonymize it and to decide whether it is sampled. /// See [`objectstore_inventory_tracker`] for the record format. /// -/// Object identities remain their storage paths. Session identities are JSON tuples of -/// `("upload_session", object_storage_path, session_id)`. Every operation on a session -/// uses the same identity. Sampling decisions are consistent while the backend -/// [`sample_rate`](CostTrackerStreamConfig::sample_rate) is unchanged; configuration changes -/// can change the decision between a session write and delete. Object and session sampling -/// decisions may differ. Neither raw paths nor session tokens are emitted. -/// -/// Sessions emit `WRITE_SESSION` and `DELETE_SESSION`. High-volume upload markers retain -/// their existing identities but emit ordinary `WRITE` and `DELETE`, since their stored -/// bytes contribute to cost attribution. -/// /// Logs, counts, and swallows errors returned by the [`InventoryTracker`]. pub struct CostTrackerStream { tracker: InventoryTracker

, diff --git a/objectstore-service/src/change_stream/mod.rs b/objectstore-service/src/change_stream/mod.rs index cc3b3a04..e1eff868 100644 --- a/objectstore-service/src/change_stream/mod.rs +++ b/objectstore-service/src/change_stream/mod.rs @@ -6,9 +6,9 @@ //! //! # Resumable uploads //! -//! Upload backends report a session write after creation with exactly the advertised upload -//! length and an accounting expiration set by `UPLOAD_SESSION_TTL`, independent of object -//! expiration. Partial chunks and incomplete offset queries emit nothing. Successful publication +//! Upload backends report a session write after creation with the advertised upload length +//! plus backend metadata bytes, and an accounting expiration set by `UPLOAD_SESSION_TTL`, +//! independent of object expiration. Partial chunks and incomplete offset queries emit nothing. Successful publication //! reports the object's actual stored size and expiration, followed by a session delete. Completion //! discovered through an offset query follows the same order. Successful cancellation also //! reports a session delete. Session writes and deletes use `WRITE_SESSION` and `DELETE_SESSION` @@ -52,7 +52,7 @@ pub(crate) use factory::dummy_factory; /// How long a backend waits for reported records to be handed off during shutdown. pub const FLUSH_TIMEOUT: Duration = Duration::from_secs(2); -/// Accounting lifetime for an upload's advertised size, independent of backend cleanup. +/// Accounting lifetime for an upload's estimated stored size, independent of backend cleanup. pub(crate) const UPLOAD_SESSION_TTL: Duration = Duration::from_hours(7 * 24); /// The object or upload session whose storage changed. @@ -145,10 +145,6 @@ fn default_sample_rate() -> f64 { #[async_trait::async_trait] pub trait ChangeStream: fmt::Debug + Send + Sync + 'static { /// Reports that `target` now occupies `size` bytes. Used for new writes and overwrites. - /// - /// Upload sessions optimistically report the advertised upload length with an accounting - /// expiration set by `UPLOAD_SESSION_TTL`. Markers report their stored size and actual - /// expiration instead. fn write(&self, target: ChangeTarget<'_>, size: u64, expires_at: Option); /// Reports that `target`'s expiration moved, with its stored size unchanged. From 1c34f2abb5fc980ff4806aa4ee0884779b0131c2 Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Thu, 8 Oct 2026 16:52:34 +0200 Subject: [PATCH 06/11] ref(change-stream): Report upload markers as objects --- objectstore-service/docs/architecture.md | 5 ++--- objectstore-service/src/backend/bigtable.rs | 12 ++++------- objectstore-service/src/backend/in_memory.rs | 5 ++--- .../src/change_stream/cost_tracker.rs | 10 +++------ objectstore-service/src/change_stream/mod.rs | 21 ++----------------- 5 files changed, 13 insertions(+), 40 deletions(-) diff --git a/objectstore-service/docs/architecture.md b/objectstore-service/docs/architecture.md index 7cb2ef1b..6942118b 100644 --- a/objectstore-service/docs/architecture.md +++ b/objectstore-service/docs/architecture.md @@ -194,15 +194,14 @@ per-backend feed of three operations: unchanged. In practice this is a TTI bump. - `delete(target)`: `target` was deleted explicitly. -[`ChangeTarget`](change_stream::ChangeTarget) identifies an object, a stored upload -marker, or an upload session. +[`ChangeTarget`](change_stream::ChangeTarget) identifies an object or an upload session. 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 storage as well as for the tombstone record in high-volume storage. -For objects and markers, `size` is a count of bytes that the backend actually stores. +For objects, `size` is a count of bytes that the backend actually stores. This includes object payloads, metadata, and sometimes backend-specific overhead. For upload sessions, `size` is the final `size` of the corresponding `object` that the upload will create when the upload is finalized. diff --git a/objectstore-service/src/backend/bigtable.rs b/objectstore-service/src/backend/bigtable.rs index cc943833..f5d74d26 100644 --- a/objectstore-service/src/backend/bigtable.rs +++ b/objectstore-service/src/backend/bigtable.rs @@ -54,7 +54,7 @@ use crate::backend::common::{ Tombstone, }; use crate::change_stream::{ - ChangeStream, ChangeStreamFactory, ChangeTarget, CostTrackerStreamConfig, flush_change_stream, + ChangeStream, ChangeStreamFactory, CostTrackerStreamConfig, flush_change_stream, }; use crate::error::{Error, ErrorKind, Result, ResultExt as _}; use crate::gcp_auth::PrefetchingTokenProvider; @@ -1090,11 +1090,8 @@ impl HighVolumeBackend for BigTableBackend { 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( - ChangeTarget::upload_marker(revision), - size, - Some(time_expires), - ); + self.change_stream + .write(revision.into(), size, Some(time_expires)); Ok(()) } @@ -1133,8 +1130,7 @@ impl HighVolumeBackend for BigTableBackend { ) .await?; if deleted { - self.change_stream - .delete(ChangeTarget::upload_marker(revision)); + self.change_stream.delete(revision.into()); } Ok(deleted) } diff --git a/objectstore-service/src/backend/in_memory.rs b/objectstore-service/src/backend/in_memory.rs index 01aa53b9..00dc31f7 100644 --- a/objectstore-service/src/backend/in_memory.rs +++ b/objectstore-service/src/backend/in_memory.rs @@ -384,7 +384,7 @@ impl HighVolumeBackend for InMemoryBackend { .unwrap() .insert(revision.clone(), time_expires); self.change_stream - .write(ChangeTarget::upload_marker(revision), 1, Some(time_expires)); + .write(revision.into(), 1, Some(time_expires)); Ok(()) } @@ -409,8 +409,7 @@ impl HighVolumeBackend for InMemoryBackend { .remove(revision) .is_some_and(|expiry| expiry >= access_time); if deleted { - self.change_stream - .delete(ChangeTarget::upload_marker(revision)); + self.change_stream.delete(revision.into()); } Ok(deleted) } diff --git a/objectstore-service/src/change_stream/cost_tracker.rs b/objectstore-service/src/change_stream/cost_tracker.rs index d45f27d5..859ba925 100644 --- a/objectstore-service/src/change_stream/cost_tracker.rs +++ b/objectstore-service/src/change_stream/cost_tracker.rs @@ -73,14 +73,12 @@ impl<'a> ChangeTarget<'a> { fn inventory_identity(self) -> (&'a ObjectId, String) { let (object_id, session_id) = match self { Self::Object(id) => return (id, id.as_storage_path().to_string()), - Self::UploadMarker(revision) => (revision, revision.key.as_str()), Self::UploadSession { object_id, session_id, } => (object_id, session_id), }; - // Preserve the established identity for both sessions and markers. The operation - // type distinguishes session estimates from stored markers on the wire. + // Keep session identities separate from objects and concurrent uploads to the same object. let key = serde_json::to_string(&("upload_session", object_id.as_storage_path(), session_id)) .expect("session identity is serializable"); @@ -97,7 +95,7 @@ where fn write(&self, target: ChangeTarget<'_>, size: u64, expires_at: Option) { let (id, key) = target.inventory_identity(); let write = match target { - ChangeTarget::Object(_) | ChangeTarget::UploadMarker(_) => InventoryTracker::write, + ChangeTarget::Object(_) => InventoryTracker::write, ChangeTarget::UploadSession { .. } => InventoryTracker::write_session, }; let result = write( @@ -129,9 +127,7 @@ where fn delete(&self, target: ChangeTarget<'_>) { let (id, key) = target.inventory_identity(); let result = match target { - ChangeTarget::Object(_) | ChangeTarget::UploadMarker(_) => { - self.tracker.delete(&key, id.usecase(), SystemTime::now()) - } + ChangeTarget::Object(_) => self.tracker.delete(&key, id.usecase(), SystemTime::now()), ChangeTarget::UploadSession { .. } => { self.tracker .delete_session(&key, id.usecase(), SystemTime::now()) diff --git a/objectstore-service/src/change_stream/mod.rs b/objectstore-service/src/change_stream/mod.rs index e1eff868..545245eb 100644 --- a/objectstore-service/src/change_stream/mod.rs +++ b/objectstore-service/src/change_stream/mod.rs @@ -15,13 +15,6 @@ //! on the inventory wire, allowing consumers to exclude their estimates from cost attribution. //! The accounting deadline does not change backend cleanup behavior. //! -//! High-volume markers use ordinary `WRITE`/`DELETE` operations to report their stored size -//! and actual deadline, set by -//! `backend::tiered::RESUMABLE_UPLOAD_TTL`. Only a successful -//! conditional marker deletion emits a delete. Tiered deletes its marker to claim the upload -//! before finalization, so that marker delete precedes the long-term object write and session -//! delete. Failed operations leave accounting intact until successful cleanup or expiration. -//! //! Behind the `storage-cogs` feature. Without it every backend gets a [`NoopStream`] and //! the transport is left out of the binary. @@ -58,14 +51,11 @@ pub(crate) const UPLOAD_SESSION_TTL: Duration = Duration::from_hours(7 * 24); /// The object or upload session whose storage changed. /// /// Session identities are separate from objects, including concurrent uploads to the same object. -/// Upload backends use their opaque backend token as `session_id`; high-volume markers use -/// [`ChangeTarget::upload_marker`]. +/// Upload backends use their opaque backend token as `session_id`. #[derive(Clone, Copy, Debug)] pub enum ChangeTarget<'a> { - /// A published object, including a high-volume redirect tombstone. + /// An object stored by a backend. Object(&'a ObjectId), - /// A stored high-volume upload marker, accounted for as an ordinary object. - UploadMarker(&'a ObjectId), /// An in-progress upload, identified by its backend session token. UploadSession { /// Object identity supplying the usecase and scopes. @@ -75,13 +65,6 @@ pub enum ChangeTarget<'a> { }, } -impl<'a> ChangeTarget<'a> { - /// Identifies the high-volume marker for an upload's unique revision. - pub fn upload_marker(revision: &'a ObjectId) -> Self { - Self::UploadMarker(revision) - } -} - impl<'a> From<&'a ObjectId> for ChangeTarget<'a> { fn from(id: &'a ObjectId) -> Self { Self::Object(id) From 9e269e33ad0984844cfb34b73bd97ef8c3825cfe Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Thu, 8 Oct 2026 16:55:55 +0200 Subject: [PATCH 07/11] improve --- objectstore-inventory-tracker/src/lib.rs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/objectstore-inventory-tracker/src/lib.rs b/objectstore-inventory-tracker/src/lib.rs index 9ab79877..1a30be63 100644 --- a/objectstore-inventory-tracker/src/lib.rs +++ b/objectstore-inventory-tracker/src/lib.rs @@ -8,7 +8,8 @@ //! of the identifier passed in by the caller. //! - `op_type`: `WRITE`, `UPDATE`, or `DELETE` for stored objects; `WRITE_SESSION` or //! `DELETE_SESSION` for upload sessions. -//! - `size`: stored bytes (including metadata), or the estimated upload-session size. +//! - `size`: stored bytes (including metadata). For sessions, this is optimistically the +//! final size that the object will have when the upload is completed. //! - `expiration_time`: a timestamp (unixtime microseconds) describing when the record is //! meant to be deleted. //! From 33b8090403702d020f51d9606648babfb931b8d1 Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Thu, 8 Oct 2026 17:05:33 +0200 Subject: [PATCH 08/11] ref(backends): Move upload session TTL to common module --- objectstore-service/src/backend/common.rs | 3 +++ objectstore-service/src/backend/gcs.rs | 5 ++--- objectstore-service/src/backend/in_memory.rs | 6 ++---- objectstore-service/src/backend/local_fs.rs | 5 ++--- objectstore-service/src/change_stream/mod.rs | 6 ++---- 5 files changed, 11 insertions(+), 14 deletions(-) diff --git a/objectstore-service/src/backend/common.rs b/objectstore-service/src/backend/common.rs index 560c095a..ff0ff501 100644 --- a/objectstore-service/src/backend/common.rs +++ b/objectstore-service/src/backend/common.rs @@ -25,6 +25,9 @@ use crate::stream::{ClientStream, PayloadStream}; /// This intentionally has a "sentry" prefix so that it can easily be traced back to us. pub const USER_AGENT: &str = concat!("sentry-objectstore/", env!("CARGO_PKG_VERSION")); +/// Accounting lifetime for an upload's estimated stored size, independent of backend cleanup. +pub(crate) const UPLOAD_SESSION_TTL: Duration = Duration::from_hours(7 * 24); + /// Backend response for put operations. pub type PutResponse = (); /// Backend response for get operations. diff --git a/objectstore-service/src/backend/gcs.rs b/objectstore-service/src/backend/gcs.rs index 48cbd294..00c0fa22 100644 --- a/objectstore-service/src/backend/gcs.rs +++ b/objectstore-service/src/backend/gcs.rs @@ -20,12 +20,11 @@ use serde::{Deserialize, Serialize}; use crate::backend::common::{ self, Backend, DeleteResponse, ExpiryUpdate, GetResponse, MetadataResponse, - MultipartUploadBackend, PutResponse, SetExpiryResponse, + MultipartUploadBackend, PutResponse, SetExpiryResponse, UPLOAD_SESSION_TTL, }; use crate::backend::extensions::{ReqwestResultExt, ResponseExt, SendTraced}; use crate::change_stream::{ - ChangeStream, ChangeStreamFactory, ChangeTarget, CostTrackerStreamConfig, UPLOAD_SESSION_TTL, - flush_change_stream, + ChangeStream, ChangeStreamFactory, ChangeTarget, CostTrackerStreamConfig, flush_change_stream, }; use crate::error::{Error, ErrorKind, Result, ResultExt as _}; use crate::gcp_auth::PrefetchingTokenProvider; diff --git a/objectstore-service/src/backend/in_memory.rs b/objectstore-service/src/backend/in_memory.rs index 00dc31f7..7f195db4 100644 --- a/objectstore-service/src/backend/in_memory.rs +++ b/objectstore-service/src/backend/in_memory.rs @@ -20,11 +20,9 @@ use objectstore_types::metadata::Metadata; use crate::backend::common::{ self, DeleteResponse, ExpiryUpdate, GetResponse, HighVolumeBackend, MultipartUploadBackend, PutResponse, SetExpiryResponse, TieredGet, TieredMetadata, TieredUpdate, TieredWrite, - Tombstone, -}; -use crate::change_stream::{ - ChangeStream, ChangeTarget, NoopStream, UPLOAD_SESSION_TTL, flush_change_stream, + Tombstone, UPLOAD_SESSION_TTL, }; +use crate::change_stream::{ChangeStream, ChangeTarget, NoopStream, flush_change_stream}; use crate::error::{Error, ErrorKind, Result}; use crate::id::ObjectId; use crate::multipart::{ diff --git a/objectstore-service/src/backend/local_fs.rs b/objectstore-service/src/backend/local_fs.rs index 2fe52565..81cc5fd3 100644 --- a/objectstore-service/src/backend/local_fs.rs +++ b/objectstore-service/src/backend/local_fs.rs @@ -39,11 +39,10 @@ use uuid::Uuid; use crate::backend::common::{ self, Backend, DeleteResponse, ExpiryUpdate, GetResponse, MultipartUploadBackend, PutResponse, - SetExpiryResponse, + SetExpiryResponse, UPLOAD_SESSION_TTL, }; use crate::change_stream::{ - ChangeStream, ChangeStreamFactory, ChangeTarget, CostTrackerStreamConfig, UPLOAD_SESSION_TTL, - flush_change_stream, + ChangeStream, ChangeStreamFactory, ChangeTarget, CostTrackerStreamConfig, flush_change_stream, }; use crate::error::{Error, ErrorKind, Result, ResultExt as _}; use crate::id::ObjectId; diff --git a/objectstore-service/src/change_stream/mod.rs b/objectstore-service/src/change_stream/mod.rs index 545245eb..7a6703f6 100644 --- a/objectstore-service/src/change_stream/mod.rs +++ b/objectstore-service/src/change_stream/mod.rs @@ -7,7 +7,8 @@ //! # Resumable uploads //! //! Upload backends report a session write after creation with the advertised upload length -//! plus backend metadata bytes, and an accounting expiration set by `UPLOAD_SESSION_TTL`, +//! plus backend metadata bytes, and an accounting expiration set by +//! [`UPLOAD_SESSION_TTL`](crate::backend::common::UPLOAD_SESSION_TTL), //! independent of object expiration. Partial chunks and incomplete offset queries emit nothing. Successful publication //! reports the object's actual stored size and expiration, followed by a session delete. Completion //! discovered through an offset query follows the same order. Successful cancellation also @@ -45,9 +46,6 @@ pub(crate) use factory::dummy_factory; /// How long a backend waits for reported records to be handed off during shutdown. pub const FLUSH_TIMEOUT: Duration = Duration::from_secs(2); -/// Accounting lifetime for an upload's estimated stored size, independent of backend cleanup. -pub(crate) const UPLOAD_SESSION_TTL: Duration = Duration::from_hours(7 * 24); - /// The object or upload session whose storage changed. /// /// Session identities are separate from objects, including concurrent uploads to the same object. From 5f8be69859359f77f3f65d874837103295483ec3 Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Thu, 8 Oct 2026 17:07:02 +0200 Subject: [PATCH 09/11] ref(change-stream): Rename UploadSession variant to Session --- objectstore-service/src/backend/gcs.rs | 2 +- objectstore-service/src/backend/in_memory.rs | 2 +- objectstore-service/src/backend/local_fs.rs | 2 +- objectstore-service/src/change_stream/cost_tracker.rs | 8 ++++---- objectstore-service/src/change_stream/mod.rs | 4 ++-- 5 files changed, 9 insertions(+), 9 deletions(-) diff --git a/objectstore-service/src/backend/gcs.rs b/objectstore-service/src/backend/gcs.rs index 00c0fa22..2ae25885 100644 --- a/objectstore-service/src/backend/gcs.rs +++ b/objectstore-service/src/backend/gcs.rs @@ -1215,7 +1215,7 @@ impl Backend for GcsBackend { })?; let token = String::from(session_uri); self.change_stream.write( - ChangeTarget::UploadSession { + ChangeTarget::Session { object_id: id, session_id: &token, }, diff --git a/objectstore-service/src/backend/in_memory.rs b/objectstore-service/src/backend/in_memory.rs index 7f195db4..8db4cd6d 100644 --- a/objectstore-service/src/backend/in_memory.rs +++ b/objectstore-service/src/backend/in_memory.rs @@ -271,7 +271,7 @@ impl super::common::Backend for InMemoryBackend { Arc::new(tokio::sync::Mutex::new(Some(upload))), ); self.change_stream.write( - ChangeTarget::UploadSession { + ChangeTarget::Session { object_id: id, session_id: &token, }, diff --git a/objectstore-service/src/backend/local_fs.rs b/objectstore-service/src/backend/local_fs.rs index 81cc5fd3..4ab205ce 100644 --- a/objectstore-service/src/backend/local_fs.rs +++ b/objectstore-service/src/backend/local_fs.rs @@ -327,7 +327,7 @@ impl Backend for LocalFsBackend { let metadata_size = UploadFile::create(&path, metadata).await?; let token = upload_id.to_string(); self.change_stream.write( - ChangeTarget::UploadSession { + ChangeTarget::Session { object_id: id, session_id: &token, }, diff --git a/objectstore-service/src/change_stream/cost_tracker.rs b/objectstore-service/src/change_stream/cost_tracker.rs index 859ba925..b92347b8 100644 --- a/objectstore-service/src/change_stream/cost_tracker.rs +++ b/objectstore-service/src/change_stream/cost_tracker.rs @@ -73,7 +73,7 @@ impl<'a> ChangeTarget<'a> { fn inventory_identity(self) -> (&'a ObjectId, String) { let (object_id, session_id) = match self { Self::Object(id) => return (id, id.as_storage_path().to_string()), - Self::UploadSession { + Self::Session { object_id, session_id, } => (object_id, session_id), @@ -96,7 +96,7 @@ where let (id, key) = target.inventory_identity(); let write = match target { ChangeTarget::Object(_) => InventoryTracker::write, - ChangeTarget::UploadSession { .. } => InventoryTracker::write_session, + ChangeTarget::Session { .. } => InventoryTracker::write_session, }; let result = write( &self.tracker, @@ -128,7 +128,7 @@ where let (id, key) = target.inventory_identity(); let result = match target { ChangeTarget::Object(_) => self.tracker.delete(&key, id.usecase(), SystemTime::now()), - ChangeTarget::UploadSession { .. } => { + ChangeTarget::Session { .. } => { self.tracker .delete_session(&key, id.usecase(), SystemTime::now()) } @@ -289,7 +289,7 @@ mod tests { let mut record_ids = std::collections::HashSet::new(); for i in 0..32 { let session_id = format!("session-{i}"); - let target = ChangeTarget::UploadSession { + let target = ChangeTarget::Session { object_id: &id, session_id: &session_id, }; diff --git a/objectstore-service/src/change_stream/mod.rs b/objectstore-service/src/change_stream/mod.rs index 7a6703f6..94cce414 100644 --- a/objectstore-service/src/change_stream/mod.rs +++ b/objectstore-service/src/change_stream/mod.rs @@ -55,7 +55,7 @@ pub enum ChangeTarget<'a> { /// An object stored by a backend. Object(&'a ObjectId), /// An in-progress upload, identified by its backend session token. - UploadSession { + Session { /// Object identity supplying the usecase and scopes. object_id: &'a ObjectId, /// Stable identity for this particular upload. @@ -71,7 +71,7 @@ impl<'a> From<&'a ObjectId> for ChangeTarget<'a> { impl<'a> From<&'a Session> for ChangeTarget<'a> { fn from(session: &'a Session) -> Self { - Self::UploadSession { + Self::Session { object_id: &session.object_id, session_id: &session.backend_token, } From 788b37670206df5767b984b494b3178ac673a28d Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Thu, 8 Oct 2026 17:11:57 +0200 Subject: [PATCH 10/11] ref(cogs): Simplify session identity encoding --- objectstore-service/src/change_stream/cost_tracker.rs | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/objectstore-service/src/change_stream/cost_tracker.rs b/objectstore-service/src/change_stream/cost_tracker.rs index b92347b8..4c8fbe21 100644 --- a/objectstore-service/src/change_stream/cost_tracker.rs +++ b/objectstore-service/src/change_stream/cost_tracker.rs @@ -79,9 +79,7 @@ impl<'a> ChangeTarget<'a> { } => (object_id, session_id), }; // Keep session identities separate from objects and concurrent uploads to the same object. - let key = - serde_json::to_string(&("upload_session", object_id.as_storage_path(), session_id)) - .expect("session identity is serializable"); + let key = format!("session:{}:{session_id}", object_id.as_storage_path()); (object_id, key) } } From e36aeed82fa410326946e8ec504cb0799d224924 Mon Sep 17 00:00:00 2001 From: lcian <17258265+lcian@users.noreply.github.com> Date: Thu, 8 Oct 2026 17:13:05 +0200 Subject: [PATCH 11/11] improve --- objectstore-inventory-tracker/src/record.rs | 2 +- objectstore-service/src/backend/common.rs | 2 +- objectstore-service/src/change_stream/mod.rs | 12 ------------ 3 files changed, 2 insertions(+), 14 deletions(-) diff --git a/objectstore-inventory-tracker/src/record.rs b/objectstore-inventory-tracker/src/record.rs index aacda9d8..7368eb45 100644 --- a/objectstore-inventory-tracker/src/record.rs +++ b/objectstore-inventory-tracker/src/record.rs @@ -1,7 +1,7 @@ //! The wire format emitted onto the inventory topic. //! //! These types mirror the `shared-resources-inventory` schema registered in -//! [sentry-kafka-schemas]. The schema sets `additionalProperties: false`, so adding a +//! [sentry-kafka-schemas]. The schema sets `additionalProperties: false`, so adding a //! field here without a corresponding schema version bump produces messages that //! consumers reject. //! diff --git a/objectstore-service/src/backend/common.rs b/objectstore-service/src/backend/common.rs index ff0ff501..c3879d78 100644 --- a/objectstore-service/src/backend/common.rs +++ b/objectstore-service/src/backend/common.rs @@ -25,7 +25,7 @@ use crate::stream::{ClientStream, PayloadStream}; /// This intentionally has a "sentry" prefix so that it can easily be traced back to us. pub const USER_AGENT: &str = concat!("sentry-objectstore/", env!("CARGO_PKG_VERSION")); -/// Accounting lifetime for an upload's estimated stored size, independent of backend cleanup. +/// Lifetime for a resumable upload session. pub(crate) const UPLOAD_SESSION_TTL: Duration = Duration::from_hours(7 * 24); /// Backend response for put operations. diff --git a/objectstore-service/src/change_stream/mod.rs b/objectstore-service/src/change_stream/mod.rs index 94cce414..36b757a0 100644 --- a/objectstore-service/src/change_stream/mod.rs +++ b/objectstore-service/src/change_stream/mod.rs @@ -4,18 +4,6 @@ //! the service describes where those records go with a [`CostTrackerConfig`], shared by //! every backend. [`ChangeStreamFactory`] pairs the two into a [`ChangeStream`]. //! -//! # Resumable uploads -//! -//! Upload backends report a session write after creation with the advertised upload length -//! plus backend metadata bytes, and an accounting expiration set by -//! [`UPLOAD_SESSION_TTL`](crate::backend::common::UPLOAD_SESSION_TTL), -//! independent of object expiration. Partial chunks and incomplete offset queries emit nothing. Successful publication -//! reports the object's actual stored size and expiration, followed by a session delete. Completion -//! discovered through an offset query follows the same order. Successful cancellation also -//! reports a session delete. Session writes and deletes use `WRITE_SESSION` and `DELETE_SESSION` -//! on the inventory wire, allowing consumers to exclude their estimates from cost attribution. -//! The accounting deadline does not change backend cleanup behavior. -//! //! Behind the `storage-cogs` feature. Without it every backend gets a [`NoopStream`] and //! the transport is left out of the binary.