diff --git a/objectstore-inventory-tracker/src/lib.rs b/objectstore-inventory-tracker/src/lib.rs index f7ed4dbe..1a30be63 100644 --- a/objectstore-inventory-tracker/src/lib.rs +++ b/objectstore-inventory-tracker/src/lib.rs @@ -6,7 +6,10 @@ //! 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. +//! - `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. //! diff --git a/objectstore-inventory-tracker/src/record.rs b/objectstore-inventory-tracker/src/record.rs index bb3378c7..7368eb45 100644 --- a/objectstore-inventory-tracker/src/record.rs +++ b/objectstore-inventory-tracker/src/record.rs @@ -11,9 +11,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 +23,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 +65,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 +168,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 f195610f..6942118b 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,17 +148,17 @@ 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 @@ -166,14 +166,17 @@ 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. 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 +184,27 @@ 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 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. -`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, `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. 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..f5d74d26 100644 --- a/objectstore-service/src/backend/bigtable.rs +++ b/objectstore-service/src/backend/bigtable.rs @@ -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,11 @@ impl HighVolumeBackend for BigTableBackend { timestamp_micros: time_expires.as_micros() as i64, value: vec![1], }))]; - self.mutate( - revision.as_upload_path().to_string().into_bytes(), - mutations, - "create_upload_marker", - ) - .await?; + let path = revision.as_upload_path().to_string().into_bytes(); + let size = row_size(&path, &mutations); + self.mutate(path, mutations, "create_upload_marker").await?; + self.change_stream + .write(revision.into(), size, Some(time_expires)); Ok(()) } @@ -1118,13 +1118,21 @@ impl HighVolumeBackend for BigTableBackend { revision: &ObjectId, access_time: Timestamp, ) -> Result { - self.check_and_mutate( - revision.as_upload_path().to_string().into_bytes(), - MutatePredicate::Include(live_row_filter(column_filter(COLUMN_UPLOAD), access_time)), - vec![delete_row_mutation()], - "delete_upload_marker", - ) - .await + let deleted = self + .check_and_mutate( + revision.as_upload_path().to_string().into_bytes(), + MutatePredicate::Include(live_row_filter( + column_filter(COLUMN_UPLOAD), + access_time, + )), + vec![delete_row_mutation()], + "delete_upload_marker", + ) + .await?; + if deleted { + self.change_stream.delete(revision.into()); + } + Ok(deleted) } #[tracing::instrument(level = "debug", fields(?id), skip_all)] @@ -1151,7 +1159,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 +1377,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 +1408,7 @@ impl HighVolumeBackend for BigTableBackend { .await?; if deleted { - self.change_stream.delete(id); + self.change_stream.delete(id.into()); return Ok(None); } @@ -1487,10 +1496,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/common.rs b/objectstore-service/src/backend/common.rs index 560c095a..c3879d78 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")); +/// Lifetime for a resumable upload session. +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 7b870f92..2ae25885 100644 --- a/objectstore-service/src/backend/gcs.rs +++ b/objectstore-service/src/backend/gcs.rs @@ -20,11 +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, CostTrackerStreamConfig, flush_change_stream, + ChangeStream, ChangeStreamFactory, ChangeTarget, CostTrackerStreamConfig, flush_change_stream, }; use crate::error::{Error, ErrorKind, Result, ResultExt as _}; use crate::gcp_auth::PrefetchingTokenProvider; @@ -752,7 +752,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 +1100,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 +1140,7 @@ impl Backend for GcsBackend { .await?; if deleted { - self.change_stream.delete(id); + self.change_stream.delete(id.into()); } Ok(()) @@ -1155,7 +1155,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", )?; @@ -1211,7 +1213,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::Session { + object_id: id, + session_id: &token, + }, + upload_length.get().saturating_add(metadata_size), + Some(Timestamp::now() + UPLOAD_SESSION_TTL), + ); + Ok(Some(token)) } #[tracing::instrument(level = "debug", fields(?session, offset, content_length), skip_all)] @@ -1270,6 +1281,7 @@ impl Backend for GcsBackend { object.metadata_size(), expires_at, ); + self.change_stream.delete(session.into()); } Ok(progress.into()) } @@ -1305,6 +1317,7 @@ impl Backend for GcsBackend { object.metadata_size(), expires_at, ); + self.change_stream.delete(session.into()); } Ok(progress.into()) }) @@ -1350,7 +1363,9 @@ impl Backend for GcsBackend { } } }) - .await + .await?; + self.change_stream.delete(session.into()); + Ok(()) } async fn join(&self) { @@ -3243,12 +3258,30 @@ 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(); 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 + GcsObject::from_metadata(&metadata).metadata_size()) + ); + 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 +3295,33 @@ 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::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); + assert_eq!(records[0].size, records[1].size); + 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::DeleteSession); + assert_eq!(records[2].op_type, OpType::DeleteSession); + assert!(records.iter().all(|r| r.record_id == records[0].record_id)); Ok(()) } @@ -3306,10 +3359,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::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); + 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..8db4cd6d 100644 --- a/objectstore-service/src/backend/in_memory.rs +++ b/objectstore-service/src/backend/in_memory.rs @@ -20,9 +20,9 @@ use objectstore_types::metadata::Metadata; use crate::backend::common::{ self, DeleteResponse, ExpiryUpdate, GetResponse, HighVolumeBackend, MultipartUploadBackend, PutResponse, SetExpiryResponse, TieredGet, TieredMetadata, TieredUpdate, TieredWrite, - Tombstone, + Tombstone, UPLOAD_SESSION_TTL, }; -use crate::change_stream::{ChangeStream, NoopStream, flush_change_stream}; +use crate::change_stream::{ChangeStream, ChangeTarget, NoopStream, flush_change_stream}; use crate::error::{Error, ErrorKind, Result}; use crate::id::ObjectId; use crate::multipart::{ @@ -181,7 +181,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 +238,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 +250,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 +259,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 +270,16 @@ impl super::common::Backend for InMemoryBackend { (id.clone(), token.clone()), Arc::new(tokio::sync::Mutex::new(Some(upload))), ); + self.change_stream.write( + ChangeTarget::Session { + object_id: id, + session_id: &token, + }, + upload_length + .get() + .saturating_add(json_len(metadata) as u64), + 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(revision.into(), 1, Some(time_expires)); Ok(()) } @@ -386,12 +400,16 @@ impl HighVolumeBackend for InMemoryBackend { revision: &ObjectId, access_time: Timestamp, ) -> Result { - Ok(self + let deleted = self .upload_markers .lock() .unwrap() .remove(revision) - .is_some_and(|expiry| expiry >= access_time)) + .is_some_and(|expiry| expiry >= access_time); + if deleted { + self.change_stream.delete(revision.into()); + } + Ok(deleted) } async fn put_non_tombstone( @@ -414,7 +432,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 +493,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 +522,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 +548,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 +747,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 +1136,67 @@ 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 + json_len(&Metadata::default()) 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) + ); - // 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::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); + assert_eq!(records[0].size, records[1].size); + 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::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 751c64b1..4ab205ce 100644 --- a/objectstore-service/src/backend/local_fs.rs +++ b/objectstore-service/src/backend/local_fs.rs @@ -39,10 +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, CostTrackerStreamConfig, flush_change_stream, + ChangeStream, ChangeStreamFactory, ChangeTarget, CostTrackerStreamConfig, flush_change_stream, }; use crate::error::{Error, ErrorKind, Result, ResultExt as _}; use crate::id::ObjectId; @@ -192,7 +192,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 +283,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 +299,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 +316,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); @@ -324,8 +324,17 @@ impl Backend for LocalFsBackend { ErrorKind::BackendFailure, "creating local-fs object directory", )?; - UploadFile::create(&path, metadata).await?; - Ok(Some(upload_id.to_string())) + let metadata_size = UploadFile::create(&path, metadata).await?; + let token = upload_id.to_string(); + self.change_stream.write( + ChangeTarget::Session { + object_id: id, + session_id: &token, + }, + upload_length.get().saturating_add(metadata_size), + 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 +386,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 +404,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 +754,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( @@ -821,7 +834,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); @@ -829,11 +842,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 { @@ -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::WriteSession, OpType::Write, OpType::DeleteSession] + ); + 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,33 @@ 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); + 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() + ..=(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 +2345,28 @@ 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::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); + 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) + ); + + 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::DeleteSession); + 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..02fd2a28 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,47 @@ mod tests { .await? .unwrap(); assert_eq!(stream::read_to_vec(body).await?, payload); + #[cfg(feature = "storage-cogs")] + { + use objectstore_inventory_tracker::OpType::{ + Delete, DeleteSession, Write, WriteSession, + }; + let records = producer.records(); + assert_eq!( + records + .iter() + .map(|r| (r.shared_resource_id.as_str(), r.op_type)) + .collect::>(), + [ + ("gcs_objectstore", WriteSession), + ("bigtable_objectstore", Write), + ("bigtable_objectstore", Delete), + ("gcs_objectstore", Write), + ("gcs_objectstore", DeleteSession), + ("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::>(), + [WriteSession, Write, Delete, DeleteSession] + ); + 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..4c8fbe21 100644 --- a/objectstore-service/src/change_stream/cost_tracker.rs +++ b/objectstore-service/src/change_stream/cost_tracker.rs @@ -6,15 +6,16 @@ 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. +/// See [`objectstore_inventory_tracker`] for the record format. /// /// Logs, counts, and swallows errors returned by the [`InventoryTracker`]. pub struct CostTrackerStream { @@ -67,15 +68,37 @@ 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) { + let (object_id, session_id) = match self { + Self::Object(id) => return (id, id.as_storage_path().to_string()), + Self::Session { + object_id, + session_id, + } => (object_id, session_id), + }; + // Keep session identities separate from objects and concurrent uploads to the same object. + let key = format!("session:{}:{session_id}", object_id.as_storage_path()); + (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) { - let result = self.tracker.write( - &id.as_storage_path().to_string(), + fn write(&self, target: ChangeTarget<'_>, size: u64, expires_at: Option) { + let (id, key) = target.inventory_identity(); + let write = match target { + ChangeTarget::Object(_) => InventoryTracker::write, + ChangeTarget::Session { .. } => InventoryTracker::write_session, + }; + let result = write( + &self.tracker, + &key, id.usecase(), size, SystemTime::now(), @@ -86,9 +109,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 +122,15 @@ 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 = match target { + ChangeTarget::Object(_) => self.tracker.delete(&key, id.usecase(), SystemTime::now()), + ChangeTarget::Session { .. } => { + self.tracker + .delete_session(&key, id.usecase(), SystemTime::now()) + } + }; self.swallow("delete", result); } @@ -140,7 +167,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 +183,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 +199,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 +232,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 +247,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 +267,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 +279,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::Session { + 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::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)); + 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..36b757a0 100644 --- a/objectstore-service/src/change_stream/mod.rs +++ b/objectstore-service/src/change_stream/mod.rs @@ -16,6 +16,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 +34,38 @@ 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); +/// 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`. +#[derive(Clone, Copy, Debug)] +pub enum ChangeTarget<'a> { + /// An object stored by a backend. + Object(&'a ObjectId), + /// An in-progress upload, identified by its backend session token. + Session { + /// Object identity supplying the usecase and scopes. + object_id: &'a ObjectId, + /// Stable identity for this particular upload. + session_id: &'a str, + }, +} + +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::Session { + 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 +108,19 @@ 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. + 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 +146,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]