From d5463936b134ce88c6a05d7e2b38fb032064e4c7 Mon Sep 17 00:00:00 2001 From: Joris Bayer Date: Tue, 25 Aug 2026 15:48:04 +0200 Subject: [PATCH 1/7] Revert "fix(upload): Feature-flag multipart (#6240)" This reverts commit 032a216cff72320af3ca71c321a2ba5f1f16af53. --- relay-dynamic-config/src/feature.rs | 5 ----- relay-server/src/endpoints/common.rs | 1 - relay-server/src/endpoints/upload.rs | 10 +--------- relay-server/src/services/objectstore.rs | 13 ++++++------- relay-server/src/services/upload.rs | 12 +++--------- tests/integration/test_upload.py | 7 +++---- 6 files changed, 13 insertions(+), 35 deletions(-) diff --git a/relay-dynamic-config/src/feature.rs b/relay-dynamic-config/src/feature.rs index 8ad4f8422e0..073af36c20c 100644 --- a/relay-dynamic-config/src/feature.rs +++ b/relay-dynamic-config/src/feature.rs @@ -92,11 +92,6 @@ pub enum Feature { /// Stream minidumps to objectstore. #[serde(rename = "projects:relay-minidump-uploads")] MinidumpUploads, - /// Use objectstore multipart for upload requests. - /// - /// See . - #[serde(rename = "projects:relay-upload-multipart")] - UploadMultipart, /// Split an NVIDIA GPU crash dump (`.nv-gpudmp`) off a minidump upload into its /// own event. #[serde(rename = "organizations:gpu-crash-symbolication")] diff --git a/relay-server/src/endpoints/common.rs b/relay-server/src/endpoints/common.rs index 599b6b07e40..9168ab63ce9 100644 --- a/relay-server/src/endpoints/common.rs +++ b/relay-server/src/endpoints/common.rs @@ -596,7 +596,6 @@ where project: project.clone(), length: None, attachment_type: item.attachment_type(), - multipart: false, }) .await .map_err(|_| BadStoreRequest::UploadFailed)? diff --git a/relay-server/src/endpoints/upload.rs b/relay-server/src/endpoints/upload.rs index 7bceb81c234..fe5744738dd 100644 --- a/relay-server/src/endpoints/upload.rs +++ b/relay-server/src/endpoints/upload.rs @@ -185,18 +185,12 @@ async fn handle_post( StatusCode::SERVICE_UNAVAILABLE })?; - let multipart = match project.state() { - ProjectState::Enabled(p) => p.has_feature(Feature::UploadMultipart), - _ => false, - }; - relay_log::trace!("Checking request"); let project_context = validate_and_limit(&state, meta, &headers, project).await?; // Unconditionally create the upload location: relay_log::trace!("Creating upload location"); - - let result = create(&state, project_context, &headers, multipart).await; + let result = create(&state, project_context, &headers).await; let location = result.inspect_err(|e| { relay_log::warn!(error = e as &dyn std::error::Error, "create failed"); })?; @@ -319,7 +313,6 @@ async fn create( state: &ServiceState, project: ProjectContext, headers: &tus::Headers, - multipart: bool, ) -> Result, Error> { let location = state .upload() @@ -327,7 +320,6 @@ async fn create( project, length: headers.upload_length, attachment_type: headers.metadata.map(|m| m.attachment_type), - multipart, }) .await??; diff --git a/relay-server/src/services/objectstore.rs b/relay-server/src/services/objectstore.rs index 7d1b736e445..e59d9a45109 100644 --- a/relay-server/src/services/objectstore.rs +++ b/relay-server/src/services/objectstore.rs @@ -50,7 +50,7 @@ pub enum Objectstore { TraceAttachment(Managed), EventAttachment(Managed), RawProfile(Managed), - Create(CreateMultipart, Sender>), + Create(Create, Sender>), Stream(Stream, Sender>), } @@ -137,7 +137,7 @@ impl MessageKind { } /// A request to create a new objectstore multipart upload. -pub struct CreateMultipart { +pub struct Create { /// The sentry org. pub organization_id: OrganizationId, /// The sentry project. @@ -148,10 +148,10 @@ pub struct CreateMultipart { pub retention: u16, } -impl FromMessage for Objectstore { +impl FromMessage for Objectstore { type Response = AsyncResponse>; - fn from_message(message: CreateMultipart, sender: Sender>) -> Self { + fn from_message(message: Create, sender: Sender>) -> Self { Self::Create(message, sender) } } @@ -792,8 +792,8 @@ impl ObjectstoreServiceInner { Ok(Some(stored_key)) } - async fn handle_create(&self, create: CreateMultipart) -> Result { - let CreateMultipart { + async fn handle_create(&self, create: Create) -> Result { + let Create { organization_id, project_id, key, @@ -811,7 +811,6 @@ impl ObjectstoreServiceInner { .send() .await?; debug_assert_eq!(&key, multipart_upload.key()); - let upload_id = multipart_upload.upload_id(); Ok(UploadRef { diff --git a/relay-server/src/services/upload.rs b/relay-server/src/services/upload.rs index d676b28688b..f33ed3bca33 100644 --- a/relay-server/src/services/upload.rs +++ b/relay-server/src/services/upload.rs @@ -138,8 +138,6 @@ pub struct Create { pub length: Option, /// The attachment type of the upload. pub attachment_type: Option, - /// Whether multipart uploads should be used for this upload. - pub multipart: bool, } /// The type used to stream a request body. @@ -271,12 +269,10 @@ impl Service { project, length, attachment_type, - multipart, }: Create, ) -> Result, Error> { match &self.backend { Backend::Upstream { addr } => { - let _ = multipart; // upstream will check feature flag again, no need to propagate. let (request, rx) = UploadRequest::create(project, length, attachment_type); addr.send(SendRequest(request)); let response = rx.await??; @@ -297,13 +293,11 @@ impl Service { .. } = project.scoping; - let (key, upload_id) = match (multipart, length) { - // We should only create a multipart upload in objectstore if it was requested, - // and if the upload actually has data (multipart does not allow empty parts). - (false, _) | (_, Some(0)) => (key, None), + let (key, upload_id) = match length { + Some(0) => (key, None), // multipart does not allow empty uploads _ => { let UploadRef { key, upload_id } = addr - .send(objectstore::CreateMultipart { + .send(objectstore::Create { organization_id, project_id, key, diff --git a/tests/integration/test_upload.py b/tests/integration/test_upload.py index 9e70fc77337..efffa7601d1 100644 --- a/tests/integration/test_upload.py +++ b/tests/integration/test_upload.py @@ -571,7 +571,7 @@ def do_upload(): [pytest.param(False, id="no multipart"), pytest.param(True, id="with multipart")], ) def test_objectstore_retries( - mini_sentry, relay_with_processing, with_multipart, project_config + mini_sentry, relay_with_processing, project_config, with_multipart ): project_id = 42 project_key = mini_sentry.get_dsn_public_key(project_id) @@ -609,6 +609,7 @@ def test_objectstore_retries( }, data=data, ) + print(response.text) failure = mini_sentry.test_failures.get(timeout=10) expected_attempts = 1 if with_multipart else 3 # multipart cannot be retried @@ -619,12 +620,10 @@ def test_objectstore_retries( assert response.status_code == 500 -def test_objectstore_timeout(mini_sentry, relay_with_processing): +def test_objectstore_timeout(mini_sentry, relay_with_processing, project_config): mini_sentry.allow_chunked = True mini_sentry.fail_on_relay_error = False project_id = 42 - config = mini_sentry.add_full_project_config(project_id)["config"] - config.setdefault("features", []).append("projects:relay-upload-multipart") project_key = mini_sentry.get_dsn_public_key(project_id) @mini_sentry.app.route( From d70c1d7f02f962cf0d2951f28ff4944c9d14d860 Mon Sep 17 00:00:00 2001 From: Joris Bayer Date: Tue, 25 Aug 2026 15:48:39 +0200 Subject: [PATCH 2/7] Revert "fix(upload): Retry multipart uploads (#6208)" This reverts commit 28ad2786fd029b82d6cec42fae227f05373923ea. --- relay-server/src/services/objectstore.rs | 24 +++--------------------- tests/integration/test_upload.py | 11 +++++++---- 2 files changed, 10 insertions(+), 25 deletions(-) diff --git a/relay-server/src/services/objectstore.rs b/relay-server/src/services/objectstore.rs index e59d9a45109..de843fc4e16 100644 --- a/relay-server/src/services/objectstore.rs +++ b/relay-server/src/services/objectstore.rs @@ -995,27 +995,9 @@ impl ObjectstoreServiceInner { while let Some((i, chunk)) = body.next().await { let chunk = chunk?; let part_number = u32::try_from(i + 1) - .map_err(|_| objectstore_client::Error::InvalidPartNumber(u32::MAX))?; - relay_log::trace!("Part number {part_number}"); - - // NOTE: This is a retry loop within a retry loop (see caller of this function). - // if we keep the Rechunked approach we might as well remove the outer loop for streaming uploads. - let mut attempts = 0; - let part = loop { - let result = multipart_upload.put(chunk.clone(), part_number, None).await; - attempts += 1; - if attempts < self.max_attempts.get() - && matches!(&result, Err(e) if is_retryable(e)) - { - relay_log::trace!("Attempt {attempts}: Failed with {result:?}, retrying"); - tokio::time::sleep(self.retry_interval).await; - } else { - relay_log::trace!("Final attempt"); - break result; - } - }; - - parts.push(part?); + .map_err(|_| objectstore_client::Error::InvalidPartNumber(u32::MAX))?; + let part = multipart_upload.put(chunk, part_number, None).await?; + parts.push(part); } multipart_upload.complete(parts).await? }); diff --git a/tests/integration/test_upload.py b/tests/integration/test_upload.py index efffa7601d1..16b3fda437b 100644 --- a/tests/integration/test_upload.py +++ b/tests/integration/test_upload.py @@ -620,7 +620,9 @@ def test_objectstore_retries( assert response.status_code == 500 -def test_objectstore_timeout(mini_sentry, relay_with_processing, project_config): +def test_objectstore_timeout( + mini_sentry, relay_with_processing, project_config, dummy_upload +): mini_sentry.allow_chunked = True mini_sentry.fail_on_relay_error = False project_id = 42 @@ -630,6 +632,7 @@ def test_objectstore_timeout(mini_sentry, relay_with_processing, project_config) "/v1/objects:multipart/attachments//", methods=["PUT"] ) def multipart_create(**params): + print(params) return {"key": params["key"], "upload_id": "foo"}, 201 @mini_sentry.app.route( @@ -637,14 +640,14 @@ def multipart_create(**params): ) def multipart_upload(**opts): time.sleep(2) - return 204 + raise NotImplementedError relay = relay_with_processing( options={ "processing": { "objectstore": { "objectstore_url": mini_sentry.url, - "stream_timeout": 1, + "timeout": 1, } } } @@ -652,7 +655,7 @@ def multipart_upload(**opts): response = upload_something(relay, project_id, project_key) - assert response.status_code == 504 + assert response.status_code == 500 # not 504 def upload_something(relay, project_id, project_key): From c429e73bdd7edd9518dce94d6f9addb103c3496d Mon Sep 17 00:00:00 2001 From: Joris Bayer Date: Tue, 25 Aug 2026 16:24:47 +0200 Subject: [PATCH 3/7] feat(upload): Revert multipart --- Cargo.lock | 1 - Cargo.toml | 4 +- relay-server/src/services/objectstore.rs | 82 ++------- relay-server/src/utils/stream/mod.rs | 4 - relay-server/src/utils/stream/rechunked.rs | 188 --------------------- 5 files changed, 19 insertions(+), 260 deletions(-) delete mode 100644 relay-server/src/utils/stream/rechunked.rs diff --git a/Cargo.lock b/Cargo.lock index 34db2b0d0ca..27bb2c0a647 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3366,7 +3366,6 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d76659c9da42fc50c394edda51e9d69dc6b3a6a770029e50f8e267c85cbfe69f" dependencies = [ "async-compression", - "base64", "bytes", "futures-util", "infer", diff --git a/Cargo.toml b/Cargo.toml index 33e421e9ce5..e775401d49e 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -161,8 +161,8 @@ metrics = "0.24" metrics-exporter-dogstatsd = "0.9" num-traits = "0.2" num_cpus = "1" -objectstore-client = { version = "0.2.1", features = ["multipart"] } -objectstore-types = { version = "0.2.0" } +objectstore-client = { version = "0.2" } +objectstore-types = { version = "0.2" } opentelemetry-semantic-conventions = "0.31" opentelemetry-proto = { version = "0.31", default-features = false } papaya = "0.2" diff --git a/relay-server/src/services/objectstore.rs b/relay-server/src/services/objectstore.rs index de843fc4e16..4affcc10778 100644 --- a/relay-server/src/services/objectstore.rs +++ b/relay-server/src/services/objectstore.rs @@ -1,20 +1,18 @@ //! Objectstore service for uploading attachments. use std::array::TryFromSliceError; use std::fmt; -use std::num::{NonZeroU16, NonZeroUsize}; +use std::num::NonZeroU16; use std::sync::Arc; use std::time::Duration; -use async_compression::tokio::bufread::ZstdEncoder; use bytes::Bytes; use futures::StreamExt; use http::StatusCode; use objectstore_client::{ - Client, Compression, ExpirationPolicy, SecretKey as SigningKey, Session, TokenGenerator, - UploadId, Usecase, + Client, ExpirationPolicy, SecretKey as SigningKey, Session, TokenGenerator, Usecase, }; -use objectstore_types::multipart::InvalidUploadId; +use objectstore_types::multipart::{InvalidUploadId, UploadId}; use relay_base_schema::organization::OrganizationId; use relay_base_schema::project::ProjectId; use relay_config::ObjectstoreServiceConfig; @@ -23,7 +21,6 @@ use relay_system::{ Addr, AsyncResponse, FromMessage, Interface, LoadShed, NoResponse, Sender, SimpleService, }; use sentry_protos::snuba::v1::TraceItem; -use tokio_util::io::{ReaderStream, StreamReader}; use crate::constants::DEFAULT_ATTACHMENT_RETENTION; use crate::envelope::{ContentType, Item, ItemType}; @@ -35,15 +32,10 @@ use crate::services::store::{ }; use crate::services::upload::ByteStream; use crate::statsd::{RelayCounters, RelayTimers}; -use crate::utils::{ - BoundedStream, MeteredStream, Rechunk, RetryableStream, TakeOnce, find_error_source, -}; +use crate::utils::{BoundedStream, MeteredStream, RetryableStream, TakeOnce, find_error_source}; use super::outcome::Outcome; -/// Size of an individual request to objectstore. -const CHUNK_SIZE: NonZeroUsize = NonZeroUsize::new(5 * 1024 * 1024).unwrap(); - /// Messages that the objectstore service can handle. pub enum Objectstore { Event(Managed>), @@ -337,7 +329,7 @@ pub struct UploadRef { /// They key of the file (chosen by relay). pub key: String, /// The ID of the multipart upload session (chosen by objectstore). - /// `None` if the upload is not multipart. + /// `None` if the upload is not a resumable session. pub upload_id: Option, } @@ -801,21 +793,12 @@ impl ObjectstoreServiceInner { } = create; let session = self.session(&self.event_attachments, organization_id, project_id)?; - let multipart_upload = session - .initiate_multipart_upload() - .expiration_policy(ExpirationPolicy::TimeToLive(Duration::from_hours( - u64::from(retention) * 24, - ))) - .key(&key) - .compression(Compression::Zstd) // make explicit because parts need to be manually compressed. - .send() - .await?; - debug_assert_eq!(&key, multipart_upload.key()); - let upload_id = multipart_upload.upload_id(); + // This is intentionally a stub. Once Objectstore implements resumable uploads, + // create an upload session here. Ok(UploadRef { key, - upload_id: Some(upload_id.clone()), + upload_id: None, }) } @@ -961,48 +944,17 @@ impl ObjectstoreServiceInner { upload_ref, retention, } => { - let UploadRef { key, upload_id } = upload_ref; - let Some(upload_id) = upload_id else { - // No upload ID: simple upload in a single request. - let request = session.put_stream(body.boxed()).key(key); - let response = request - .expiration_policy(ExpirationPolicy::TimeToLive(Duration::from_hours( - u64::from(retention) * 24, - ))) - .send() - .await?; - return Ok(ObjectstoreKey(response.key)); - }; + let UploadRef { key, upload_id: _ } = upload_ref; - let multipart_upload = - session.resume_multipart_upload(key, upload_id.to_string())?; + let request = session.put_stream(body.boxed()).key(key); + let response = request + .expiration_policy(ExpirationPolicy::TimeToLive(Duration::from_hours( + u64::from(retention) * 24, + ))) + .send() + .await?; - let body = ReaderStream::new(ZstdEncoder::new(StreamReader::new(body))); - - // Unfortunately, MinIO has the limitation that the length of a multipart request - // has to be known. Therefore, we need to materialize the stream into concrete - // chunks of bytes and send each chunk as an individual request. - let chunks = Rechunk::new(body, CHUNK_SIZE); - let mut body = chunks.enumerate(); - - let result = relay_statsd::metric!( - timer(RelayTimers::AttachmentUploadDuration), - type = kind.as_str(), - { - let mut parts = vec![]; - // NOTE: Once every upload is a multipart upload, we can remove `RetryableStream` - // because streams will never be effectively retried. - while let Some((i, chunk)) = body.next().await { - let chunk = chunk?; - let part_number = u32::try_from(i + 1) - .map_err(|_| objectstore_client::Error::InvalidPartNumber(u32::MAX))?; - let part = multipart_upload.put(chunk, part_number, None).await?; - parts.push(part); - } - multipart_upload.complete(parts).await? - }); - - Ok(ObjectstoreKey(result)) + Ok(ObjectstoreKey(response.key)) } } } diff --git a/relay-server/src/utils/stream/mod.rs b/relay-server/src/utils/stream/mod.rs index a487744edc4..c4956e29006 100644 --- a/relay-server/src/utils/stream/mod.rs +++ b/relay-server/src/utils/stream/mod.rs @@ -1,13 +1,9 @@ mod bounded; mod metered; mod peek; -#[cfg(any(feature = "processing", test))] -mod rechunked; mod retryable; pub use bounded::*; pub use metered::*; pub use peek::*; -#[cfg(feature = "processing")] -pub use rechunked::*; pub use retryable::*; diff --git a/relay-server/src/utils/stream/rechunked.rs b/relay-server/src/utils/stream/rechunked.rs deleted file mode 100644 index b3527a4d741..00000000000 --- a/relay-server/src/utils/stream/rechunked.rs +++ /dev/null @@ -1,188 +0,0 @@ -use std::num::NonZeroUsize; -use std::pin::Pin; -use std::task::{Context, Poll}; - -use bytes::{BufMut, Bytes, BytesMut}; -use futures::{Stream, StreamExt}; - -/// A stream adapter that emits chunks with a fixed size. -/// -/// All emitted [`Bytes`] have `chunk_size` bytes except for the final chunk, which may be smaller. -pub struct Rechunk { - inner: S, - chunk_size: usize, - buffer: BytesMut, - /// State of the stream: - /// - `None`: not done. - /// - `Some(Some(e))`: need to flush an error. - /// - `Some(None)`: completely done. - done: Option>, -} - -impl Rechunk { - /// Creates a new stream adapter that emits chunks of `chunk_size`. - pub fn new(inner: S, chunk_size: NonZeroUsize) -> Self { - Self { - inner, - chunk_size: chunk_size.get(), - buffer: BytesMut::new(), - done: None, - } - } -} - -impl Rechunk -where - S: Stream> + Unpin, -{ - fn flush_one(&mut self) -> Poll>> { - let chunk = self.buffer.split_to(self.chunk_size.min(self.buffer.len())); - if chunk.is_empty() { - if let Some(done) = &mut self.done { - if let Some(error) = done.take() { - // Flush the error. Will be done on next poll. - return Poll::Ready(Some(Err(error))); - } else { - return Poll::Ready(None); - } - }; - return Poll::Pending; - } - - Poll::Ready(Some(Ok(chunk.freeze()))) - } -} - -impl Stream for Rechunk -where - S: Stream> + Unpin, - E: Unpin, -{ - type Item = Result; - - fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { - let this = self.get_mut(); - - // While there are full chunks, flush them. - if this.buffer.len() >= this.chunk_size { - return this.flush_one(); - } - - // Flush final chunk or error. - if this.done.is_some() { - return this.flush_one(); - } - - while this.buffer.len() < this.chunk_size { - match this.inner.poll_next_unpin(cx) { - Poll::Ready(None) => { - this.done = Some(None); - break; - } - Poll::Ready(Some(Err(e))) => { - this.done = Some(Some(e)); - break; - } - Poll::Ready(Some(Ok(bytes))) => { - this.buffer.put(bytes); - } - Poll::Pending => return Poll::Pending, - } - } - - this.flush_one() - } -} - -#[cfg(test)] -mod tests { - use futures::stream; - - use super::*; - - fn bytes(value: &'static [u8]) -> Bytes { - Bytes::from_static(value) - } - - async fn collect_chunks( - chunks: Vec>, - chunk_size: usize, - ) -> Vec> { - Rechunk::new(stream::iter(chunks), NonZeroUsize::new(chunk_size).unwrap()) - .map(|a| a) - .collect::>() - .await - } - - #[tokio::test] - async fn test_splits_large_chunks() { - let chunks = collect_chunks(vec![Ok(bytes(b"abcdefg"))], 3).await; - - assert_eq!( - chunks, - vec![Ok(bytes(b"abc")), Ok(bytes(b"def")), Ok(bytes(b"g"))] - ); - } - - #[tokio::test] - async fn test_combines_small_chunks() { - let chunks = collect_chunks( - vec![Ok(bytes(b"ab")), Ok(bytes(b"c")), Ok(bytes(b"defg"))], - 3, - ) - .await; - - assert_eq!( - chunks, - vec![Ok(bytes(b"abc")), Ok(bytes(b"def")), Ok(bytes(b"g"))] - ); - } - - #[tokio::test] - async fn test_yields_exact_chunks_unchanged() { - let chunks = collect_chunks(vec![Ok(bytes(b"ab")), Ok(bytes(b"cd"))], 2).await; - - assert_eq!(chunks, vec![Ok(bytes(b"ab")), Ok(bytes(b"cd"))]); - } - - #[tokio::test] - async fn test_skips_empty_chunks() { - let chunks = collect_chunks( - vec![ - Ok(bytes(b"")), - Ok(bytes(b"ab")), - Ok(bytes(b"")), - Ok(bytes(b"cde")), - Ok(bytes(b"")), - ], - 2, - ) - .await; - - assert_eq!( - chunks, - vec![Ok(bytes(b"ab")), Ok(bytes(b"cd")), Ok(bytes(b"e"))] - ); - } - - #[tokio::test] - async fn test_empty_stream() { - let chunks = collect_chunks(Vec::new(), 2).await; - - assert!(chunks.is_empty()); - } - - #[tokio::test] - async fn test_propagates_errors() { - let chunks = collect_chunks(vec![Ok(bytes(b"ab")), Err("failed")], 2).await; - - assert_eq!(chunks, vec![Ok(bytes(b"ab")), Err("failed")]); - } - - #[tokio::test] - async fn test_forwards_incomplete_chunk_on_error() { - let chunks = collect_chunks(vec![Ok(bytes(b"a")), Err("failed")], 2).await; - - assert_eq!(chunks, vec![Ok(bytes(b"a")), Err("failed")]); - } -} From 29515e1c31eba671238b988d0c41a448229d544c Mon Sep 17 00:00:00 2001 From: Joris Bayer Date: Tue, 25 Aug 2026 16:56:23 +0200 Subject: [PATCH 4/7] linit --- relay-server/src/services/objectstore.rs | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/relay-server/src/services/objectstore.rs b/relay-server/src/services/objectstore.rs index 59d3d679367..cf988ec6e09 100644 --- a/relay-server/src/services/objectstore.rs +++ b/relay-server/src/services/objectstore.rs @@ -23,7 +23,6 @@ use relay_system::{ Addr, AsyncResponse, FromMessage, Interface, LoadShed, NoResponse, Sender, SimpleService, }; use sentry_protos::snuba::v1::{AnyValue, TraceItem, any_value}; -use tokio_util::io::{ReaderStream, StreamReader}; use crate::constants::DEFAULT_ATTACHMENT_RETENTION; use crate::envelope::{ContentType, Item, ItemType}; @@ -819,9 +818,9 @@ impl ObjectstoreServiceInner { organization_id, project_id, key, - retention, + retention: _, } = create; - let session = self.session(&self.event_attachments, organization_id, project_id)?; + let _session = self.session(&self.event_attachments, organization_id, project_id)?; // This is intentionally a stub. Once Objectstore implements resumable uploads, // create an upload session here. From 8e59ec1b6c015759fa0162c6059122ad6e2b7cb1 Mon Sep 17 00:00:00 2001 From: Joris Bayer Date: Thu, 27 Aug 2026 12:07:57 +0200 Subject: [PATCH 5/7] test --- tests/integration/test_upload.py | 21 ++------------------- 1 file changed, 2 insertions(+), 19 deletions(-) diff --git a/tests/integration/test_upload.py b/tests/integration/test_upload.py index 16b3fda437b..0cbe09678ec 100644 --- a/tests/integration/test_upload.py +++ b/tests/integration/test_upload.py @@ -628,16 +628,7 @@ def test_objectstore_timeout( project_id = 42 project_key = mini_sentry.get_dsn_public_key(project_id) - @mini_sentry.app.route( - "/v1/objects:multipart/attachments//", methods=["PUT"] - ) - def multipart_create(**params): - print(params) - return {"key": params["key"], "upload_id": "foo"}, 201 - - @mini_sentry.app.route( - "/v1/objects:multipart:parts/attachments//", methods=["PUT"] - ) + @mini_sentry.app.route("/v1/objects/attachments//", methods=["PUT"]) def multipart_upload(**opts): time.sleep(2) raise NotImplementedError @@ -682,18 +673,10 @@ def upload_something(relay, project_id, project_key): ) -@pytest.mark.parametrize( - "with_multipart", - [pytest.param(False, id="no multipart"), pytest.param(True, id="with multipart")], -) -def test_objectstore_retention( - mini_sentry, relay_with_processing, objectstore, with_multipart -): +def test_objectstore_retention(mini_sentry, relay_with_processing, objectstore): project_id = 42 config = mini_sentry.add_full_project_config(project_id)["config"] config["eventRetention"] = 20 - if with_multipart: - config.setdefault("features", []).append("projects:relay-upload-multipart") project_key = mini_sentry.get_dsn_public_key(project_id) relay = relay_with_processing() From 5daa4f1139426a4c12c734b63d50019bbd7a5c50 Mon Sep 17 00:00:00 2001 From: Joris Bayer Date: Thu, 27 Aug 2026 15:27:39 +0200 Subject: [PATCH 6/7] test --- tests/integration/test_upload.py | 17 +++-------------- 1 file changed, 3 insertions(+), 14 deletions(-) diff --git a/tests/integration/test_upload.py b/tests/integration/test_upload.py index 0cbe09678ec..7682683a022 100644 --- a/tests/integration/test_upload.py +++ b/tests/integration/test_upload.py @@ -566,13 +566,7 @@ def do_upload(): }, r.text -@pytest.mark.parametrize( - "with_multipart", - [pytest.param(False, id="no multipart"), pytest.param(True, id="with multipart")], -) -def test_objectstore_retries( - mini_sentry, relay_with_processing, project_config, with_multipart -): +def test_objectstore_retries(mini_sentry, relay_with_processing, project_config): project_id = 42 project_key = mini_sentry.get_dsn_public_key(project_id) @@ -590,9 +584,6 @@ def test_objectstore_retries( location = f"/api/{project_id}/upload/019cdc82ed6c7761ba21fd34b86481c2/" sep = "?" - if with_multipart: - location += "?upload_id=my_upload_id" - sep = "&" signature = SecretKey.parse(relay.secret_key).sign(location.encode()) signed_location = ( f"{location}{sep}sentry_key={project_key}&upload_signature={signature}" @@ -612,10 +603,8 @@ def test_objectstore_retries( print(response.text) failure = mini_sentry.test_failures.get(timeout=10) - expected_attempts = 1 if with_multipart else 3 # multipart cannot be retried - assert ( - f"failed to upload 1 attachment(s) to objectstore in {expected_attempts} attempt(s)" - in str(failure) + assert "failed to upload 1 attachment(s) to objectstore in 3 attempt(s)" in str( + failure ) assert response.status_code == 500 From 2fcca3b550596535d2ec305c38fb518cbd0f73c5 Mon Sep 17 00:00:00 2001 From: Joris Bayer Date: Fri, 28 Aug 2026 10:53:01 +0200 Subject: [PATCH 7/7] test: restore stream_timeout --- tests/integration/test_upload.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/tests/integration/test_upload.py b/tests/integration/test_upload.py index 7682683a022..50b1b0aa5eb 100644 --- a/tests/integration/test_upload.py +++ b/tests/integration/test_upload.py @@ -618,7 +618,7 @@ def test_objectstore_timeout( project_key = mini_sentry.get_dsn_public_key(project_id) @mini_sentry.app.route("/v1/objects/attachments//", methods=["PUT"]) - def multipart_upload(**opts): + def slow_upload(**opts): time.sleep(2) raise NotImplementedError @@ -627,7 +627,7 @@ def multipart_upload(**opts): "processing": { "objectstore": { "objectstore_url": mini_sentry.url, - "timeout": 1, + "stream_timeout": 1, } } } @@ -635,7 +635,7 @@ def multipart_upload(**opts): response = upload_something(relay, project_id, project_key) - assert response.status_code == 500 # not 504 + assert response.status_code == 504 def upload_something(relay, project_id, project_key):