Skip to content
5 changes: 4 additions & 1 deletion objectstore-inventory-tracker/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
//!
Expand Down
14 changes: 11 additions & 3 deletions objectstore-inventory-tracker/src/record.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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.
Expand Down Expand Up @@ -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<u64>,

Expand Down Expand Up @@ -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));
}
Expand Down
93 changes: 91 additions & 2 deletions objectstore-inventory-tracker/src/tracker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -153,6 +153,58 @@ impl<P: Producer> InventoryTracker<P> {
expiration_time: Option<SystemTime>,
organization_id: Option<u64>,
project_id: Option<u64>,
) -> 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<SystemTime>,
organization_id: Option<u64>,
project_id: Option<u64>,
) -> 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<SystemTime>,
organization_id: Option<u64>,
project_id: Option<u64>,
) -> Result<(), P::Error> {
let Some(record_id) = self.sample(storage_key) else {
return Ok(());
Expand All @@ -161,7 +213,7 @@ impl<P: Producer> InventoryTracker<P> {
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,
Expand Down Expand Up @@ -211,6 +263,28 @@ impl<P: Producer> InventoryTracker<P> {
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(());
Expand All @@ -219,7 +293,7 @@ impl<P: Producer> InventoryTracker<P> {
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,
Expand Down Expand Up @@ -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);
Expand All @@ -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]
Expand Down
33 changes: 20 additions & 13 deletions objectstore-service/docs/architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -148,56 +148,63 @@ 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.

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`.

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
Expand Down
49 changes: 29 additions & 20 deletions objectstore-service/src/backend/bigtable.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(())
}
Expand Down Expand Up @@ -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(())
}
Expand All @@ -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(())
}

Expand All @@ -1118,13 +1118,21 @@ impl HighVolumeBackend for BigTableBackend {
revision: &ObjectId,
access_time: Timestamp,
) -> Result<bool> {
self.check_and_mutate(
revision.as_upload_path().to_string().into_bytes(),
MutatePredicate::Include(live_row_filter(column_filter(COLUMN_UPLOAD), access_time)),
vec![delete_row_mutation()],
"delete_upload_marker",
)
.await
let deleted = self
.check_and_mutate(
revision.as_upload_path().to_string().into_bytes(),
MutatePredicate::Include(live_row_filter(
column_filter(COLUMN_UPLOAD),
access_time,
)),
vec![delete_row_mutation()],
"delete_upload_marker",
)
.await?;
if deleted {
self.change_stream.delete(revision.into());
}
Ok(deleted)
}

#[tracing::instrument(level = "debug", fields(?id), skip_all)]
Expand 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);
}

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

Expand Down Expand Up @@ -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)
Expand Down
3 changes: 3 additions & 0 deletions objectstore-service/src/backend/common.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
Loading
Loading