diff --git a/objectstore-server/src/config.rs b/objectstore-server/src/config.rs index f9edacfa..fb23eab1 100644 --- a/objectstore-server/src/config.rs +++ b/objectstore-server/src/config.rs @@ -683,6 +683,7 @@ impl Default for Config { storage: StorageConfig::FileSystem(FileSystemConfig { path: PathBuf::from("data"), + cogs: None, }), storage_cogs: None, diff --git a/objectstore-service/docs/architecture.md b/objectstore-service/docs/architecture.md index 36564890..3d1d0f79 100644 --- a/objectstore-service/docs/architecture.md +++ b/objectstore-service/docs/architecture.md @@ -146,16 +146,11 @@ 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, TTI bump, and delete it performs on -stored objects to a [`ChangeStream`](change_stream::ChangeStream). 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 queried to break down each backend's storage utilization -by `app_feature`. Note that [`NoopStream`](change_stream::NoopStream) is used -unless the backend's config includes a -[`CostTrackerStreamConfig`](change_stream::CostTrackerStreamConfig) and the -service has a usable transport for it, and unless the `storage-cogs` feature is -compiled in at all. +Storage attribution is derived from the [change stream](#change-streams) each +backend publishes. To turn a 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 queried to break down each backend's storage +utilization by `app_feature`. Each row in the inventory table has an anonymized hash of an `ObjectId` as well as the row's size, expiry, Sentry org/project, `app_feature`, and relevant @@ -164,11 +159,10 @@ 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. -`ChangeStream` is not aware of any automatic garbage collection that backends -may perform. Expired objects must be filtered out when querying the inventory -table. +Because the change stream does not observe automatic garbage collection, expired +objects must be filtered out when querying the inventory table. -Under the hood, `CostTrackerStream` uses +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 @@ -179,6 +173,48 @@ sampling rate. When aggregating, divide each row's value by its `sample_rate`. See also: [`objectstore_inventory_tracker`] documentation. +# Change Streams + +Every backend publishes the changes it makes to the objects it stores as a +[`ChangeStream`](change_stream::ChangeStream). It is a fire-and-forget, +per-backend feed of three operations: + +- `write(id, size, expires_at)`: `id` now occupies `size` bytes. Used for both + new objects and overwrites. +- `update(id, expires_at)`: `id`'s expiration moved while its stored size is + unchanged. In practice this is a TTI bump. +- `delete(id)`: `id` was deleted explicitly. + +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. + +Decorators such as [`CountingBackend`](backend::counting::CountingBackend) and +[`TieredStorage`](backend::tiered::TieredStorage) don't publish change streams +of their own; only leaf backends that actually own bytes do. + +Automatic garbage collection is invisible to the change stream. Downstream +consumers of the stream need to consider the `expires_at` field on messages. + +## `ChangeStream` implementation guidance + +While the [`ChangeStream`](change_stream::ChangeStream) trait is abstract, that +abstraction is not surfaced in service configuration. For instance, the +[storage COGS change stream](#storage-cogs) is configured with a service-wide +[`CostTrackerConfig`](change_stream::CostTrackerConfig) and per-backend +[`CostTrackerStreamConfig`](change_stream::CostTrackerStreamConfig)s. These +configurations are connected in [`ChangeStreamFactory`](change_stream::ChangeStreamFactory) +to build a [`CostTrackerStream`](change_stream::CostTrackerStream). + +New `ChangeStream` implementations may follow the same pattern: +- per-backend configuration for per-backend IDs or configuration +- service-wide configuration for a stream sink +- glue code in and around `ChangeStreamFactory` + # Metadata and Payload Every object consists of structured **metadata** and a binary **payload**. diff --git a/objectstore-service/src/backend/in_memory.rs b/objectstore-service/src/backend/in_memory.rs index 93e74cfc..b8d88a88 100644 --- a/objectstore-service/src/backend/in_memory.rs +++ b/objectstore-service/src/backend/in_memory.rs @@ -19,6 +19,7 @@ use super::common::{ DeleteResponse, GetResponse, HighVolumeBackend, MultipartUploadBackend, PutResponse, TieredGet, TieredMetadata, TieredWrite, Tombstone, }; +use crate::change_stream::{ChangeStream, NoopStream, flush_change_stream}; use crate::error::{Error, Result}; use crate::id::ObjectId; use crate::multipart::{ @@ -34,6 +35,20 @@ enum StoreEntry { Tombstone(Tombstone), } +impl StoreEntry { + /// Number of bytes an entry occupies: its payload plus its serialized metadata. + pub(crate) fn stored_size(&self) -> usize { + match self { + StoreEntry::Object(metadata, payload) => json_len(metadata) + payload.len(), + StoreEntry::Tombstone(tombstone) => { + // A tombstone carries no payload, only its redirect and expiry. + tombstone.target.as_storage_path().to_string().len() + + json_len(&tombstone.expiration_policy) + } + } + } +} + type Store = HashMap; #[derive(Clone, Debug)] @@ -61,6 +76,7 @@ pub struct InMemoryBackend { name: &'static str, store: Arc>, multipart_store: Arc>, + change_stream: Arc, } impl InMemoryBackend { @@ -70,9 +86,16 @@ impl InMemoryBackend { name, store: Arc::new(Mutex::new(HashMap::new())), multipart_store: Arc::new(Mutex::new(HashMap::new())), + change_stream: Arc::new(NoopStream), } } + /// Publishes this backend's changes to `change_stream`. + pub fn with_change_stream(mut self, change_stream: Arc) -> Self { + self.change_stream = change_stream; + self + } + /// Returns the stored entry for `id`, for direct inspection in tests. pub fn get(&self, id: &ObjectId) -> Entry { match self.store.lock().unwrap().get(id).cloned() { @@ -117,10 +140,11 @@ impl super::common::Backend for InMemoryBackend { stream: ClientStream, ) -> Result { let bytes: BytesMut = stream.try_collect().await?; - self.store.lock().unwrap().insert( - id.clone(), - StoreEntry::Object(metadata.clone(), bytes.freeze()), - ); + let entry = StoreEntry::Object(metadata.clone(), bytes.freeze()); + let size = entry.stored_size(); + self.store.lock().unwrap().insert(id.clone(), entry); + self.change_stream + .write(id, size as u64, metadata.time_expires); Ok(()) } @@ -153,9 +177,15 @@ impl super::common::Backend for InMemoryBackend { } async fn delete_object(&self, id: &ObjectId) -> Result { - self.store.lock().unwrap().remove(id); + if self.store.lock().unwrap().remove(id).is_some() { + self.change_stream.delete(id); + } Ok(()) } + + async fn join(&self) { + flush_change_stream(&self.change_stream).await; + } } #[async_trait::async_trait] @@ -173,7 +203,11 @@ impl HighVolumeBackend for InMemoryBackend { let mut metadata = metadata.clone(); metadata.size = Some(payload.len()); - store.insert(id.clone(), StoreEntry::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); Ok(None) } @@ -220,7 +254,9 @@ impl HighVolumeBackend for InMemoryBackend { return Ok(Some(tombstone)); } - store.remove(id); + if store.remove(id).is_some() { + self.change_stream.delete(id); + } Ok(None) } @@ -239,13 +275,26 @@ impl HighVolumeBackend for InMemoryBackend { if matches_current { match write { TieredWrite::Tombstone(tombstone) => { - store.insert(id.clone(), StoreEntry::Tombstone(tombstone)); + let expires_at = tombstone + .expiration_policy + .expires_in() + .map(|ttl| SystemTime::now() + ttl); + 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); } TieredWrite::Object(metadata, payload) => { - store.insert(id.clone(), StoreEntry::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); } TieredWrite::Delete => { - store.remove(id); + if store.remove(id).is_some() { + self.change_stream.delete(id); + } } } } @@ -374,7 +423,7 @@ impl MultipartUploadBackend for InMemoryBackend { // Validate and assemble while holding the multipart lock, but don't // remove the upload yet — a failed validation must leave it intact so // the client can retry. - let assembled = { + let (metadata, payload) = { let store = self.multipart_store.lock().unwrap(); let upload = store .get(&key) @@ -416,10 +465,11 @@ impl MultipartUploadBackend for InMemoryBackend { (metadata, payload.freeze()) }; - self.store - .lock() - .unwrap() - .insert(id.clone(), StoreEntry::Object(assembled.0, assembled.1)); + let expires_at = metadata.time_expires; + 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.multipart_store.lock().unwrap().remove(&key); @@ -427,6 +477,11 @@ impl MultipartUploadBackend for InMemoryBackend { } } +/// Serialized length of `value`, or `0` if it cannot be serialized. +fn json_len(value: &T) -> usize { + serde_json::to_string(value).map_or(0, |json| json.len()) +} + /// Returns `true` if `entry` matches the expected tombstone redirect state. /// /// - `expected = None`: matches any non-tombstone (absent or inline object). @@ -498,6 +553,11 @@ mod tests { use objectstore_types::metadata::ExpirationPolicy; use objectstore_types::scope::{Scope, Scopes}; + #[cfg(feature = "storage-cogs")] + use objectstore_inventory_tracker::OpType; + #[cfg(feature = "storage-cogs")] + use objectstore_inventory_tracker::test_utils::DummyProducer; + use super::*; use crate::backend::common::Backend; use crate::id::ObjectContext; @@ -822,4 +882,48 @@ mod tests { .unwrap(); assert!(result.is_none(), "retry with correct part should succeed"); } + + #[cfg(feature = "storage-cogs")] + fn backend_with_change_stream() -> (InMemoryBackend, DummyProducer) { + use crate::change_stream::CostTrackerStreamConfig; + + let (streams, producer) = crate::change_stream::dummy_factory(); + let change_stream = streams.build(Some(&CostTrackerStreamConfig { + shared_resource_id: "in_memory_objectstore".into(), + sample_rate: 1.0, + })); + ( + InMemoryBackend::new("test").with_change_stream(change_stream), + producer, + ) + } + + #[cfg(feature = "storage-cogs")] + #[tokio::test] + async fn change_stream_reports_writes_and_deletes() { + let (backend, producer) = backend_with_change_stream(); + let id = make_id(); + let metadata = Metadata::default(); + let payload = b"hello"; + + backend + .put_object(&id, &metadata, stream::single(payload.to_vec())) + .await + .unwrap(); + backend.delete_object(&id).await.unwrap(); + // The object is already gone, so this reports nothing. + backend.delete_object(&id).await.unwrap(); + + let records = producer.records(); + assert_eq!(records.len(), 2); + assert_eq!(records[0].op_type, OpType::Write); + assert_eq!(records[0].app_feature, "testing"); + assert_eq!( + records[0].size, + Some((json_len(&metadata) + payload.len()) as u64), + "the reported size covers metadata as well as the payload" + ); + assert!(json_len(&metadata) > 0, "metadata must contribute bytes"); + assert_eq!(records[1].op_type, OpType::Delete); + } } diff --git a/objectstore-service/src/backend/local_fs.rs b/objectstore-service/src/backend/local_fs.rs index af0b233d..6d95c2da 100644 --- a/objectstore-service/src/backend/local_fs.rs +++ b/objectstore-service/src/backend/local_fs.rs @@ -3,6 +3,7 @@ use std::io::ErrorKind; use std::path::PathBuf; use std::pin::pin; +use std::sync::Arc; use std::time::SystemTime; use futures_util::StreamExt; @@ -15,6 +16,9 @@ use tokio_util::io::{ReaderStream, StreamReader}; use crate::backend::common::{ Backend, DeleteResponse, GetResponse, MultipartUploadBackend, PutResponse, }; +use crate::change_stream::{ + ChangeStream, ChangeStreamFactory, CostTrackerStreamConfig, flush_change_stream, +}; use crate::error::{Error, Result}; use crate::id::ObjectId; use crate::multipart::{ @@ -51,18 +55,37 @@ pub struct FileSystemConfig { /// - `OS__STORAGE__TYPE=filesystem` /// - `OS__STORAGE__PATH=/path/to/storage` pub path: PathBuf, + + /// Reports what this backend stores, for per-usecase cost attribution. + /// + /// # Default + /// + /// `None`, which disables reporting for this backend. + /// + /// # Environment Variables + /// + /// - `OS__STORAGE__COGS__SHARED_RESOURCE_ID=filesystem_objectstore` + /// - `OS__STORAGE__COGS__SAMPLE_RATE=1.0` (optional) + #[serde(default, skip_serializing_if = "Option::is_none")] + pub cogs: Option, } /// Local filesystem backend for development and testing. #[derive(Debug)] pub struct LocalFsBackend { path: PathBuf, + + change_stream: Arc, } impl LocalFsBackend { /// Creates a new [`LocalFsBackend`] rooted at the directory in `config`. - pub fn new(config: FileSystemConfig) -> Self { - Self { path: config.path } + pub fn new(config: FileSystemConfig, streams: &ChangeStreamFactory) -> Self { + let FileSystemConfig { path, cogs } = config; + Self { + path, + change_stream: streams.build(cogs.as_ref()), + } } } @@ -103,7 +126,7 @@ impl Backend for LocalFsBackend { writer.write_all(metadata_json.as_bytes()).await?; writer.write_all(b"\n").await?; - tokio::io::copy(&mut reader, &mut writer) + let payload_size = tokio::io::copy(&mut reader, &mut writer) .await .map_err(|e| match stream::unpack_client_error(&e) { Some(ce) => Error::Client(ce), @@ -115,6 +138,10 @@ impl Backend for LocalFsBackend { file.sync_data().await?; drop(file); + let metadata_size = metadata_json.len() as u64 + 1; + self.change_stream + .write(id, metadata_size + payload_size, metadata.time_expires); + Ok(()) } @@ -168,14 +195,24 @@ impl Backend for LocalFsBackend { async fn delete_object(&self, id: &ObjectId) -> Result { objectstore_log::debug!("Deleting from local_fs backend"); let path = self.path.join(id.as_storage_path().to_string()); - let result = tokio::fs::remove_file(path).await; - if let Err(e) = &result - && e.kind() == ErrorKind::NotFound - { - objectstore_log::debug!("Object not found"); - } + + let result = match tokio::fs::remove_file(path).await { + Ok(()) => { + self.change_stream.delete(id); + Ok(()) + } + Err(e) if e.kind() == ErrorKind::NotFound => { + objectstore_log::debug!("Object not found"); + Ok(()) + } + Err(e) => Err(e), + }; Ok(result?) } + + async fn join(&self) { + flush_change_stream(&self.change_stream).await; + } } impl LocalFsBackend { @@ -418,13 +455,14 @@ impl MultipartUploadBackend for LocalFsBackend { writer.write_all(metadata_json.as_bytes()).await?; writer.write_all(b"\n").await?; + let mut payload_size = 0; for completed in &parts { let part_path = dir.join(format!("{}.part", completed.part_number)); let file = tokio::fs::File::open(&part_path).await?; let mut reader = BufReader::new(file); let mut header_line = String::new(); reader.read_line(&mut header_line).await?; - tokio::io::copy(&mut reader, &mut writer).await?; + payload_size += tokio::io::copy(&mut reader, &mut writer).await?; } writer.flush().await?; @@ -432,6 +470,10 @@ impl MultipartUploadBackend for LocalFsBackend { file.sync_data().await?; drop(file); + let header_size = metadata_json.len() as u64 + 1; + self.change_stream + .write(id, header_size + payload_size, metadata.time_expires); + // Clean up multipart state tokio::fs::remove_dir_all(dir).await?; @@ -449,6 +491,11 @@ mod tests { use objectstore_types::metadata::{Compression, ExpirationPolicy}; use objectstore_types::scope::{Scope, Scopes}; + #[cfg(feature = "storage-cogs")] + use objectstore_inventory_tracker::OpType; + #[cfg(feature = "storage-cogs")] + use objectstore_inventory_tracker::test_utils::DummyProducer; + use super::*; use crate::id::ObjectContext; use crate::stream; @@ -456,9 +503,13 @@ mod tests { #[tokio::test] async fn stores_metadata() { let tempdir = tempfile::tempdir().unwrap(); - let backend = LocalFsBackend::new(FileSystemConfig { - path: tempdir.path().to_path_buf(), - }); + let backend = LocalFsBackend::new( + FileSystemConfig { + path: tempdir.path().to_path_buf(), + cogs: None, + }, + &ChangeStreamFactory::default(), + ); let id = ObjectId::random(ObjectContext { usecase: "testing".into(), @@ -497,9 +548,13 @@ mod tests { #[tokio::test] async fn get_metadata_returns_metadata() { let tempdir = tempfile::tempdir().unwrap(); - let backend = LocalFsBackend::new(FileSystemConfig { - path: tempdir.path().to_path_buf(), - }); + let backend = LocalFsBackend::new( + FileSystemConfig { + path: tempdir.path().to_path_buf(), + cogs: None, + }, + &ChangeStreamFactory::default(), + ); let id = ObjectId::random(ObjectContext { usecase: "testing".into(), @@ -531,9 +586,13 @@ mod tests { #[tokio::test] async fn get_metadata_nonexistent() { let tempdir = tempfile::tempdir().unwrap(); - let backend = LocalFsBackend::new(FileSystemConfig { - path: tempdir.path().to_path_buf(), - }); + let backend = LocalFsBackend::new( + FileSystemConfig { + path: tempdir.path().to_path_buf(), + cogs: None, + }, + &ChangeStreamFactory::default(), + ); let id = ObjectId::random(ObjectContext { usecase: "testing".into(), @@ -551,11 +610,32 @@ mod tests { }) } + #[cfg(feature = "storage-cogs")] + fn make_backend_with_change_stream() -> (tempfile::TempDir, LocalFsBackend, DummyProducer) { + let tempdir = tempfile::tempdir().unwrap(); + let (streams, producer) = crate::change_stream::dummy_factory(); + let backend = LocalFsBackend::new( + FileSystemConfig { + path: tempdir.path().to_path_buf(), + cogs: Some(CostTrackerStreamConfig { + shared_resource_id: "filesystem_objectstore".into(), + sample_rate: 1.0, + }), + }, + &streams, + ); + (tempdir, backend, producer) + } + fn make_backend() -> (tempfile::TempDir, LocalFsBackend) { let tempdir = tempfile::tempdir().unwrap(); - let backend = LocalFsBackend::new(FileSystemConfig { - path: tempdir.path().to_path_buf(), - }); + let backend = LocalFsBackend::new( + FileSystemConfig { + path: tempdir.path().to_path_buf(), + cogs: None, + }, + &ChangeStreamFactory::default(), + ); (tempdir, backend) } @@ -970,4 +1050,64 @@ mod tests { .unwrap(); assert!(result.is_none(), "retry with correct part should succeed"); } + + #[cfg(feature = "storage-cogs")] + #[tokio::test] + async fn change_stream_reports_the_size_written_to_disk() { + let (tempdir, backend, producer) = make_backend_with_change_stream(); + let id = make_id(); + let metadata = Metadata { + expiration_policy: ExpirationPolicy::TimeToLive(Duration::from_hours(1)), + time_expires: Some(SystemTime::now() + Duration::from_hours(1)), + ..Default::default() + }; + let payload = b"oh hai!"; + + backend + .put_object(&id, &metadata, stream::single(payload.to_vec())) + .await + .unwrap(); + + let file = tokio::fs::read(tempdir.path().join(id.as_storage_path().to_string())) + .await + .unwrap(); + + let records = producer.records(); + assert_eq!(records.len(), 1); + assert_eq!(records[0].op_type, OpType::Write); + assert_eq!(records[0].shared_resource_id, "filesystem_objectstore"); + assert_eq!(records[0].app_feature, "testing"); + assert_eq!( + records[0].size, + Some(file.len() as u64), + "the metadata header line counts towards the reported size" + ); + assert!(records[0].expiration_time.is_some()); + } + + #[cfg(feature = "storage-cogs")] + #[tokio::test] + async fn change_stream_reports_deletes_on_success() { + let (_tempdir, backend, producer) = make_backend_with_change_stream(); + let id = make_id(); + + // Try to delete a non-existent object. Don't emit a message. + backend + .delete_object(&id) + .await + .expect("deleting a non-existent object returns Ok(())"); + assert!(producer.records().is_empty()); + + backend + .put_object(&id, &Metadata::default(), stream::single(b"hi".to_vec())) + .await + .unwrap(); + producer.clear(); + + backend.delete_object(&id).await.unwrap(); + + let records = producer.records(); + assert_eq!(records.len(), 1); + assert_eq!(records[0].op_type, OpType::Delete); + } } diff --git a/objectstore-service/src/backend/mod.rs b/objectstore-service/src/backend/mod.rs index a9c62d28..a0b6b81d 100644 --- a/objectstore-service/src/backend/mod.rs +++ b/objectstore-service/src/backend/mod.rs @@ -86,10 +86,10 @@ async fn from_leaf_config( streams: &ChangeStreamFactory, ) -> Result> { Ok(match config { - StorageConfig::FileSystem(c) => Box::new(local_fs::LocalFsBackend::new(c)), - StorageConfig::S3Compatible(c) => { - Box::new(s3_compatible::S3CompatibleBackend::without_token(c)) - } + StorageConfig::FileSystem(c) => Box::new(local_fs::LocalFsBackend::new(c, streams)), + StorageConfig::S3Compatible(c) => Box::new( + s3_compatible::S3CompatibleBackend::without_token(c, streams), + ), StorageConfig::Gcs(c) => Box::new(gcs::GcsBackend::new(c, streams).await?), StorageConfig::BigTable(c) => Box::new(bigtable::BigTableBackend::new(c, streams).await?), StorageConfig::Tiered(_) => anyhow::bail!("nested tiered storage is not supported"), @@ -142,7 +142,9 @@ async fn lt_from_config( streams: &ChangeStreamFactory, ) -> Result> { Ok(match config { - MultipartUploadStorageConfig::FileSystem(c) => Box::new(local_fs::LocalFsBackend::new(c)), + MultipartUploadStorageConfig::FileSystem(c) => { + Box::new(local_fs::LocalFsBackend::new(c, streams)) + } MultipartUploadStorageConfig::Gcs(c) => Box::new(gcs::GcsBackend::new(c, streams).await?), }) } diff --git a/objectstore-service/src/backend/s3_compatible.rs b/objectstore-service/src/backend/s3_compatible.rs index f34beec0..c982411d 100644 --- a/objectstore-service/src/backend/s3_compatible.rs +++ b/objectstore-service/src/backend/s3_compatible.rs @@ -1,5 +1,7 @@ //! S3-compatible backend with generic protocol support. +use std::sync::Arc; +use std::sync::atomic::Ordering; use std::time::SystemTime; use std::{fmt, io}; @@ -13,9 +15,12 @@ use super::extensions::{ResponseExt, SendTraced}; use crate::backend::common::{ self, Backend, DeleteResponse, GetResponse, MetadataResponse, PutResponse, }; +use crate::change_stream::{ + ChangeStream, ChangeStreamFactory, CostTrackerStreamConfig, flush_change_stream, +}; use crate::error::{Error, Result}; use crate::id::ObjectId; -use crate::stream::ClientStream; +use crate::stream::{ClientStream, counting_stream}; /// Configuration for [`S3CompatibleBackend`]. /// @@ -52,6 +57,19 @@ pub struct S3CompatibleConfig { /// /// - `OS__STORAGE__BUCKET=my-bucket` pub bucket: String, + + /// Reports what this backend stores, for per-usecase cost attribution. + /// + /// # Default + /// + /// `None`, which disables reporting for this backend. + /// + /// # Environment Variables + /// + /// - `OS__STORAGE__COGS__SHARED_RESOURCE_ID=s3_objectstore` + /// - `OS__STORAGE__COGS__SAMPLE_RATE=1.0` (optional) + #[serde(default, skip_serializing_if = "Option::is_none")] + pub cogs: Option, } /// Prefix used for custom metadata in headers for the GCS backend. @@ -100,16 +118,36 @@ pub struct S3CompatibleBackend { bucket: String, token_provider: Option, + + change_stream: Arc, } impl S3CompatibleBackend { /// Creates a new S3-compatible backend bound to the given bucket. - pub fn new(endpoint: &str, bucket: &str, token_provider: T) -> Self { + pub fn new( + config: S3CompatibleConfig, + token_provider: T, + streams: &ChangeStreamFactory, + ) -> Self { + Self::build(config, Some(token_provider), streams) + } + + fn build( + config: S3CompatibleConfig, + token_provider: Option, + streams: &ChangeStreamFactory, + ) -> Self { + let S3CompatibleConfig { + endpoint, + bucket, + cogs, + } = config; Self { client: common::reqwest_client(), - endpoint: endpoint.into(), - bucket: bucket.into(), - token_provider: Some(token_provider), + endpoint, + bucket, + token_provider, + change_stream: streams.build(cogs.as_ref()), } } @@ -119,6 +157,14 @@ impl S3CompatibleBackend { } } +/// Number of bytes the given headers occupy as stored object metadata. +fn headers_size(headers: &HeaderMap) -> u64 { + headers + .iter() + .map(|(name, value)| name.as_str().len() as u64 + value.len() as u64) + .sum() +} + /// Wraps [`Metadata::to_headers`] with GCS-specific concerns (tombstone + custom-time). fn metadata_to_gcs_headers( metadata: &Metadata, @@ -262,6 +308,7 @@ where let mut bumped = metadata.clone(); bumped.time_expires = Some(new_expire_at); self.update_metadata(id, &bumped).await?; + self.change_stream.update(id, Some(new_expire_at)); } Ok(Some((metadata, content_range, response))) @@ -302,13 +349,8 @@ impl fmt::Debug for S3CompatibleBackend { impl S3CompatibleBackend { /// Creates a new S3-compatible backend that sends unauthenticated requests. - pub fn without_token(config: S3CompatibleConfig) -> Self { - Self { - client: common::reqwest_client(), - endpoint: config.endpoint, - bucket: config.bucket, - token_provider: None, - } + pub fn without_token(config: S3CompatibleConfig, streams: &ChangeStreamFactory) -> Self { + Self::build(config, None, streams) } } @@ -326,10 +368,16 @@ impl Backend for S3CompatibleBackend { stream: ClientStream, ) -> Result { objectstore_log::debug!("Writing to s3_compatible backend"); + let headers = metadata_to_gcs_headers(metadata, GCS_CUSTOM_PREFIX)?; + let metadata_size = headers_size(&headers); + + // A successful PUT does not report the stored size back, so count what we send. + let (payload_size, counted) = counting_stream(stream); + self.request(Method::PUT, self.object_url(id)) .await? - .headers(metadata_to_gcs_headers(metadata, GCS_CUSTOM_PREFIX)?) - .body(Body::wrap_stream(stream)) + .headers(headers) + .body(Body::wrap_stream(counted)) .send_traced() .await .check_error("S3: failed to put object") @@ -337,6 +385,12 @@ impl Backend for S3CompatibleBackend { .drain_body() .await; + self.change_stream.write( + id, + metadata_size + payload_size.load(Ordering::Relaxed), + metadata.time_expires, + ); + Ok(()) } @@ -385,9 +439,14 @@ impl Backend for S3CompatibleBackend { .await? .drain_body() .await; + self.change_stream.delete(id); Ok(()) } + + async fn join(&self) { + flush_change_stream(&self.change_stream).await; + } } #[cfg(test)] @@ -410,10 +469,14 @@ mod tests { // Refer to the readme for how to set up MinIO via devservices. fn create_test_backend() -> S3CompatibleBackend { - S3CompatibleBackend::without_token(S3CompatibleConfig { - endpoint: "http://localhost:8089".into(), - bucket: "test-bucket".into(), - }) + S3CompatibleBackend::without_token( + S3CompatibleConfig { + endpoint: "http://localhost:8089".into(), + bucket: "test-bucket".into(), + cogs: None, + }, + &ChangeStreamFactory::default(), + ) } fn make_id() -> ObjectId { @@ -477,6 +540,18 @@ mod tests { assert_eq!(roundtripped.custom, metadata.custom); } + #[test] + fn headers_size_counts_names_and_values() { + let mut headers = HeaderMap::new(); + headers.insert("x-goog-meta-a", "1".parse().unwrap()); + headers.insert("x-goog-meta-bb", "22".parse().unwrap()); + + assert_eq!( + headers_size(&headers), + ("x-goog-meta-a".len() + 1 + "x-goog-meta-bb".len() + 2) as u64 + ); + } + #[tokio::test] async fn test_get_metadata_nonexistent() -> Result<()> { let backend = create_test_backend(); diff --git a/objectstore-service/src/backend/tiered.rs b/objectstore-service/src/backend/tiered.rs index 9887fe24..e6b97698 100644 --- a/objectstore-service/src/backend/tiered.rs +++ b/objectstore-service/src/backend/tiered.rs @@ -98,12 +98,12 @@ //! persisted. use std::sync::Arc; -use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::atomic::Ordering; use std::time::{Duration, SystemTime}; use base64::Engine as _; use bytes::Bytes; -use futures_util::{Stream, StreamExt}; +use futures_util::StreamExt; use objectstore_types::metadata::Metadata; use objectstore_types::range::ByteRange; use sentry::{Hub, SentryFutureExt}; @@ -121,7 +121,7 @@ use crate::multipart::{ AbortMultipartResponse, CompleteMultipartResponse, CompletedPart, InitiateMultipartResponse, ListPartsResponse, PartNumber, UploadId, UploadPartResponse, }; -use crate::stream::{ClientStream, SizedPeek}; +use crate::stream::{ClientStream, SizedPeek, counting_stream}; /// The threshold up until which we will go to the "high volume" backend. const BACKEND_SIZE_THRESHOLD: usize = 1024 * 1024; // 1 MiB @@ -556,26 +556,6 @@ impl std::fmt::Display for BackendChoice { } } -/// Wraps a stream to count the total bytes yielded by successful chunks. -/// -/// Returns the shared counter and the wrapped stream. The counter is incremented -/// as the stream is consumed, so read it only after the stream is exhausted. -fn counting_stream(stream: S) -> (Arc, impl Stream>) -where - S: Stream>, -{ - let counter = Arc::new(AtomicU64::new(0)); - - ( - counter.clone(), - stream.inspect(move |res| { - if let Ok(chunk) = res { - counter.fetch_add(chunk.len() as u64, Ordering::Relaxed); - } - }), - ) -} - /// The multipart upload state for TieredStorage. #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] struct TieredUploadId { diff --git a/objectstore-service/src/stream.rs b/objectstore-service/src/stream.rs index efaf9c43..f3e16ea7 100644 --- a/objectstore-service/src/stream.rs +++ b/objectstore-service/src/stream.rs @@ -13,6 +13,7 @@ use std::error::Error; use std::fmt; use std::io; use std::sync::Arc; +use std::sync::atomic::{AtomicU64, Ordering}; use bytes::{Bytes, BytesMut}; use futures_util::stream::BoxStream; @@ -297,6 +298,26 @@ pub fn single( futures_util::stream::once(std::future::ready(Ok(contents.into()))).boxed() } +/// Wraps a stream to count the total bytes yielded by successful chunks. +/// +/// Returns the shared counter and the wrapped stream. The counter is incremented +/// as the stream is consumed, so read it only after the stream is exhausted. +pub fn counting_stream(stream: S) -> (Arc, impl Stream>) +where + S: Stream>, +{ + let counter = Arc::new(AtomicU64::new(0)); + + ( + counter.clone(), + stream.inspect(move |res| { + if let Ok(chunk) = res { + counter.fetch_add(chunk.len() as u64, Ordering::Relaxed); + } + }), + ) +} + /// Collects a stream of `Bytes` chunks into a `Vec`. #[cfg(test)] pub(crate) async fn read_to_vec(mut stream: S) -> crate::error::Result>