This is an automated email from the ASF dual-hosted git repository.
hubcio pushed a commit to branch fix/docs-audit-connector-sources
in repository https://gitbox.apache.org/repos/asf/iggy.git
The following commit(s) were added to
refs/heads/fix/docs-audit-connector-sources by this push:
new 448b60de5 fix(connectors): verify S3 write access before consuming
448b60de5 is described below
commit 448b60de5f474cd6c5bcc04ec30902fc631e60f9
Author: Hubert Gruszecki <[email protected]>
AuthorDate: Mon Sep 14 18:47:29 2026 +0200
fix(connectors): verify S3 write access before consuming
Deferring write validation allowed misconfigured sinks to consume
and lose messages after offsets were committed. Initiate and abort
a multipart probe before opening, checking write access without
publishing an object. Require AbortMultipartUpload for cleanup.
---
core/connectors/sinks/s3_sink/README.md | 4 +-
core/connectors/sinks/s3_sink/src/sink.rs | 284 +++++++++++++++++++++++-
core/integration/tests/connectors/s3/s3_sink.rs | 18 +-
3 files changed, 292 insertions(+), 14 deletions(-)
diff --git a/core/connectors/sinks/s3_sink/README.md
b/core/connectors/sinks/s3_sink/README.md
index fcae30327..c3bf0c1aa 100644
--- a/core/connectors/sinks/s3_sink/README.md
+++ b/core/connectors/sinks/s3_sink/README.md
@@ -111,7 +111,9 @@ The pinned `aws-creds` chain tries these sources in order:
For temporary key pairs, supply the token through the environment or shared
credentials file. This is not the AWS SDK credential chain.
-Startup validates configuration and loads credentials without writing probe
objects. Bucket access is checked by the first real upload, so missing buckets,
denied writes and endpoint failures surface then. The sink requires `PutObject`
on data and loss-marker keys; neither `ListBucket` nor `DeleteObject` is
required.
+Startup initiates a multipart upload at `<prefix>/.iggy-sink-probe` (bucket
root when the prefix is empty) and aborts it before consuming messages. No
parts are uploaded and no object is published or deleted. Both steps must
succeed: missing buckets, denied writes or failed cleanup prevent startup. The
credentials need `s3:PutObject` and `s3:AbortMultipartUpload` on the probe key,
plus `s3:PutObject` on data and loss-marker keys. Neither `ListBucket` nor
`DeleteObject` is required.
+
+Each probe step uses `max_attempts` and `retry_delay` for HTTP 408, 429 and
5xx failures. Abort also retries transport failures. Initiation does not add
retries for an ambiguous transport failure because the upload ID may be lost;
rust-s3 can still retry internally. A crash or lost response can leave an
incomplete upload, so configure an `AbortIncompleteMultipartUpload` lifecycle
rule. Failed aborts report the upload ID for cleanup.
## Output Example
diff --git a/core/connectors/sinks/s3_sink/src/sink.rs
b/core/connectors/sinks/s3_sink/src/sink.rs
index c487aff33..fcf155f60 100644
--- a/core/connectors/sinks/s3_sink/src/sink.rs
+++ b/core/connectors/sinks/s3_sink/src/sink.rs
@@ -20,13 +20,16 @@ use crate::formatter;
use crate::path::{PathContext, render_s3_key};
use crate::{BufferKey, S3Sink};
use async_trait::async_trait;
-use iggy_connector_sdk::retry::retry_backoff;
+use iggy_connector_sdk::retry::{RetryPolicy, retry_async, retry_backoff};
use iggy_connector_sdk::{ConsumedMessage, Error, MessagesMetadata, Sink,
TopicMetadata};
+use s3::error::S3Error;
use std::sync::Arc;
use std::time::Duration;
use tracing::{debug, error, info, warn};
const MAX_BACKOFF: Duration = Duration::from_secs(60);
+const PROBE_KEY: &str = ".iggy-sink-probe";
+const PROBE_CONTENT_TYPE: &str = "application/octet-stream";
struct FlushPayload {
data: Vec<u8>,
@@ -43,7 +46,9 @@ impl Sink for S3Sink {
self.validate_and_parse_config()?;
- self.bucket = Some(crate::client::create_bucket(&self.config).await?);
+ let bucket = crate::client::create_bucket(&self.config).await?;
+ self.check_write_access(&bucket).await?;
+ self.bucket = Some(bucket);
info!(
"S3 sink ID: {} opened. format={}, rotation={}, max_file_size={},
template='{}'",
@@ -181,6 +186,62 @@ impl Sink for S3Sink {
}
impl S3Sink {
+ async fn check_write_access(&self, bucket: &s3::Bucket) -> Result<(),
Error> {
+ let prefix = self
+ .config
+ .prefix
+ .as_deref()
+ .unwrap_or_default()
+ .trim_matches('/');
+ let probe_key = if prefix.is_empty() {
+ PROBE_KEY.to_string()
+ } else {
+ format!("{prefix}/{PROBE_KEY}")
+ };
+ let resolved = self.resolved();
+ let policy = RetryPolicy {
+ max_attempts: resolved.max_attempts,
+ base_delay: resolved.retry_delay,
+ max_delay: MAX_BACKOFF,
+ };
+
+ // Initiation checks PutObject without publishing an object. A
transport
+ // failure can lose the upload ID, so do not retry that ambiguous
result.
+ let upload = retry_async(
+ policy,
+ &format!("S3 sink ID: {} initiate write probe", self.id),
+ |error| matches!(error, S3Error::HttpFailWithBody(status, _) if
is_retriable_status(*status)),
+ || bucket.initiate_multipart_upload(&probe_key,
PROBE_CONTENT_TYPE),
+ )
+ .await
+ .map_err(|error| Error::InitError(format!(
+ "S3 bucket '{}' write probe at '{probe_key}' failed: {error}",
bucket.name
+ )))?;
+
+ if upload.upload_id.is_empty() {
+ return Err(Error::InitError(format!(
+ "S3 bucket '{}' write probe at '{probe_key}' returned an empty
upload ID",
+ bucket.name
+ )));
+ }
+
+ retry_async(
+ policy,
+ &format!("S3 sink ID: {} abort write probe", self.id),
+ |error| match error {
+ S3Error::HttpFailWithBody(status, _) =>
is_retriable_status(*status),
+ S3Error::Reqwest(_) => true,
+ _ => false,
+ },
+ || bucket.abort_upload(&probe_key, &upload.upload_id),
+ )
+ .await
+ .map_err(|error| Error::InitError(format!(
+ "S3 bucket '{}' could not abort write probe at '{probe_key}'
(upload ID '{}'); s3:AbortMultipartUpload permission is required: {error}",
+ bucket.name, upload.upload_id
+ )))
+ }
+
#[allow(clippy::too_many_arguments)]
async fn process_messages_inner(
&self,
@@ -396,7 +457,7 @@ fn is_retriable_status(status: u16) -> bool {
#[cfg(test)]
mod tests {
use iggy_connector_sdk::{Payload, Schema};
- use wiremock::matchers::{method, path};
+ use wiremock::matchers::{any, method, path, query_param};
use wiremock::{Mock, MockServer, ResponseTemplate};
use super::*;
@@ -405,6 +466,32 @@ mod tests {
S3SinkConfig,
};
+ const TEST_UPLOAD_ID: &str = "test-upload-id";
+
+ fn probe_response(probe_key: &str, upload_id: &str) -> ResponseTemplate {
+ ResponseTemplate::new(200).set_body_string(format!(
+
"<InitiateMultipartUploadResult><Bucket>test-bucket</Bucket><Key>{probe_key}</Key><UploadId>{upload_id}</UploadId></InitiateMultipartUploadResult>"
+ ))
+ }
+
+ async fn mock_write_probe(server: &MockServer, probe_key: &str) {
+ let probe_path = format!("/test-bucket/{probe_key}");
+ Mock::given(method("POST"))
+ .and(path(&probe_path))
+ .and(query_param("uploads", ""))
+ .respond_with(probe_response(probe_key, TEST_UPLOAD_ID))
+ .expect(1)
+ .mount(server)
+ .await;
+ Mock::given(method("DELETE"))
+ .and(path(&probe_path))
+ .and(query_param("uploadId", TEST_UPLOAD_ID))
+ .respond_with(ResponseTemplate::new(204))
+ .expect(1)
+ .mount(server)
+ .await;
+ }
+
fn test_config() -> S3SinkConfig {
S3SinkConfig {
bucket: "test-bucket".to_string(),
@@ -427,10 +514,51 @@ mod tests {
}
#[test]
- fn
given_put_only_permissions_when_opening_should_write_only_uploaded_data() {
+ fn
given_an_unwritable_bucket_when_opening_should_fail_before_consumption() {
+ let runtime = tokio::runtime::Runtime::new().expect("Start test
runtime");
+ runtime.block_on(async {
+ for (status, expected_attempts) in [(403, 1), (404, 1), (503, 2)] {
+ let server = MockServer::start().await;
+ Mock::given(any())
+ .respond_with(ResponseTemplate::new(status))
+ .mount(&server)
+ .await;
+ let config = S3SinkConfig {
+ endpoint: Some(server.uri()),
+ access_key_id: Some("test-access-key".into()),
+ secret_access_key: Some("test-secret-key".into()),
+ max_attempts: Some(2),
+ retry_delay: Some("1ms".to_string()),
+ ..test_config()
+ };
+ let mut sink = S3Sink::new(1, config);
+
+ let result = sink.open().await;
+
+ assert!(
+ matches!(result, Err(Error::InitError(_))),
+ "S3 status {status} must prevent startup: {result:?}"
+ );
+ assert!(
+ sink.bucket.is_none(),
+ "An unwritable sink must remain unopened"
+ );
+ let requests = server.received_requests().await.expect("record
requests");
+ assert_eq!(requests.len(), expected_attempts);
+ assert!(requests.iter().all(|request| request.method ==
"POST"));
+ let state = sink.state.lock().await;
+ assert_eq!(state.messages_received, 0);
+ assert_eq!(state.messages_lost, 0);
+ }
+ });
+ }
+
+ #[test]
+ fn
given_a_multipart_probe_when_opening_should_abort_before_uploading_data() {
let runtime = tokio::runtime::Runtime::new().expect("Start test
runtime");
runtime.block_on(async {
let server = MockServer::start().await;
+ mock_write_probe(&server, "allowed/events/.iggy-sink-probe").await;
Mock::given(method("PUT"))
.respond_with(ResponseTemplate::new(200))
.mount(&server)
@@ -448,10 +576,10 @@ mod tests {
};
let mut sink = S3Sink::new(1, config);
sink.open().await.expect("Sink should initialize");
- assert!(
- server.received_requests().await.expect("record
requests").is_empty(),
- "Opening the sink must not create probe objects or require
other S3 permissions"
- );
+ let requests = server.received_requests().await.expect("record
requests");
+ assert_eq!(requests.len(), 2, "Startup must initiate and abort the
probe");
+ assert_eq!(requests[0].method, "POST");
+ assert_eq!(requests[1].method, "DELETE");
sink.consume(
&TopicMetadata {
@@ -474,14 +602,15 @@ mod tests {
}],
)
.await
- .expect("PutObject-only access should allow data uploads");
+ .expect("Data upload should succeed after the probe is cleaned
up");
sink.close().await.expect("Sink should close");
let requests = server.received_requests().await.expect("record
requests");
- assert_eq!(requests.len(), 1, "Only the data object should be
written");
- assert_eq!(requests[0].body, b"message");
+ assert_eq!(requests.len(), 3, "Only the data object should be
written after the probe");
+ assert_eq!(requests[2].method, "PUT");
+ assert_eq!(requests[2].body, b"message");
assert_eq!(
- requests[0].url.path(),
+ requests[2].url.path(),
"/test-bucket/allowed/events/messages/00000-00000000000000000000-00000000000000000000.bin"
);
let state = sink.state.lock().await;
@@ -490,12 +619,142 @@ mod tests {
});
}
+ #[test]
+ fn
given_a_probe_prefix_when_opening_should_use_the_configured_write_scope() {
+ let runtime = tokio::runtime::Runtime::new().expect("Start test
runtime");
+ runtime.block_on(async {
+ for (prefix, probe_key) in [
+ (None, ".iggy-sink-probe"),
+ (Some(""), ".iggy-sink-probe"),
+ (Some("/"), ".iggy-sink-probe"),
+ (Some("/allowed/events/"), "allowed/events/.iggy-sink-probe"),
+ ] {
+ let server = MockServer::start().await;
+ mock_write_probe(&server, probe_key).await;
+ let config = S3SinkConfig {
+ prefix: prefix.map(str::to_string),
+ endpoint: Some(server.uri()),
+ access_key_id: Some("test-access-key".into()),
+ secret_access_key: Some("test-secret-key".into()),
+ ..test_config()
+ };
+ let mut sink = S3Sink::new(1, config);
+
+ sink.open()
+ .await
+ .expect("Probe should succeed within the configured
prefix");
+ }
+ });
+ }
+
+ #[test]
+ fn
given_transient_probe_failures_when_opening_should_retry_the_failed_step() {
+ let runtime = tokio::runtime::Runtime::new().expect("Start test
runtime");
+ runtime.block_on(async {
+ for operation in ["POST", "DELETE"] {
+ for status in [408, 429, 503] {
+ let server = MockServer::start().await;
+ Mock::given(method(operation))
+ .and(path("/test-bucket/data/.iggy-sink-probe"))
+ .respond_with(ResponseTemplate::new(status))
+ .up_to_n_times(1)
+ .expect(1)
+ .mount(&server)
+ .await;
+ mock_write_probe(&server, "data/.iggy-sink-probe").await;
+ let config = S3SinkConfig {
+ endpoint: Some(server.uri()),
+ access_key_id: Some("test-access-key".into()),
+ secret_access_key: Some("test-secret-key".into()),
+ max_attempts: Some(2),
+ retry_delay: Some("1ms".to_string()),
+ ..test_config()
+ };
+ let mut sink = S3Sink::new(1, config);
+
+ sink.open()
+ .await
+ .expect("Transient probe failure should recover");
+ }
+ }
+ });
+ }
+
+ #[test]
+ fn given_a_failed_probe_abort_when_opening_should_refuse_startup() {
+ let runtime = tokio::runtime::Runtime::new().expect("Start test
runtime");
+ runtime.block_on(async {
+ for (status, expected_attempts) in [(403, 1), (503, 2)] {
+ let server = MockServer::start().await;
+ Mock::given(method("POST"))
+ .and(path("/test-bucket/data/.iggy-sink-probe"))
+ .and(query_param("uploads", ""))
+ .respond_with(probe_response("data/.iggy-sink-probe",
TEST_UPLOAD_ID))
+ .expect(1)
+ .mount(&server)
+ .await;
+ Mock::given(method("DELETE"))
+ .and(path("/test-bucket/data/.iggy-sink-probe"))
+ .and(query_param("uploadId", TEST_UPLOAD_ID))
+ .respond_with(ResponseTemplate::new(status))
+ .expect(expected_attempts)
+ .mount(&server)
+ .await;
+ let config = S3SinkConfig {
+ endpoint: Some(server.uri()),
+ access_key_id: Some("test-access-key".into()),
+ secret_access_key: Some("test-secret-key".into()),
+ max_attempts: Some(2),
+ retry_delay: Some("1ms".to_string()),
+ ..test_config()
+ };
+ let mut sink = S3Sink::new(1, config);
+
+ let result = sink.open().await;
+
+ assert!(matches!(result, Err(Error::InitError(_))),
"{result:?}");
+ assert!(sink.bucket.is_none(), "Cleanup must succeed before
opening");
+ }
+ });
+ }
+
+ #[test]
+ fn
given_an_empty_upload_id_when_probing_should_fail_without_a_delete_request() {
+ let runtime = tokio::runtime::Runtime::new().expect("Start test
runtime");
+ runtime.block_on(async {
+ let server = MockServer::start().await;
+ Mock::given(method("POST"))
+ .respond_with(probe_response("data/.iggy-sink-probe", ""))
+ .mount(&server)
+ .await;
+ let config = S3SinkConfig {
+ endpoint: Some(server.uri()),
+ access_key_id: Some("test-access-key".into()),
+ secret_access_key: Some("test-secret-key".into()),
+ ..test_config()
+ };
+ let mut sink = S3Sink::new(1, config);
+
+ let result = sink.open().await;
+
+ assert!(matches!(result, Err(Error::InitError(_))), "{result:?}");
+ let requests = server.received_requests().await.expect("record
requests");
+ assert_eq!(
+ requests.len(),
+ 1,
+ "An abort must identify the multipart upload"
+ );
+ assert_eq!(requests[0].method, "POST");
+ });
+ }
+
#[test]
fn given_upload_status_when_writing_should_retry_only_transient_failures()
{
let runtime = tokio::runtime::Runtime::new().expect("Start test
runtime");
runtime.block_on(async {
for status in [200, 403, 404, 408, 429, 503] {
let server = MockServer::start().await;
+ mock_write_probe(&server, "data/.iggy-sink-probe").await;
Mock::given(method("PUT"))
.and(path("/test-bucket/data/message.bin"))
.respond_with(ResponseTemplate::new(status))
@@ -546,6 +805,7 @@ mod tests {
runtime.block_on(async {
for previously_buffered in [0, 1] {
let server = MockServer::start().await;
+ mock_write_probe(&server, "data/.iggy-sink-probe").await;
let config = S3SinkConfig {
endpoint: Some(server.uri()),
access_key_id: Some("test-access-key".into()),
diff --git a/core/integration/tests/connectors/s3/s3_sink.rs
b/core/integration/tests/connectors/s3/s3_sink.rs
index 22258aba2..a3352d06e 100644
--- a/core/integration/tests/connectors/s3/s3_sink.rs
+++ b/core/integration/tests/connectors/s3/s3_sink.rs
@@ -21,7 +21,7 @@ use bytes::Bytes;
use iggy::prelude::{IggyMessage, Partitioning};
use iggy_common::Identifier;
use iggy_common::MessageClient;
-use iggy_connector_sdk::api::SinkInfoResponse;
+use iggy_connector_sdk::api::{ConnectorStatus, SinkInfoResponse};
use integration::harness::seeds;
use integration::iggy_harness;
use reqwest::Client;
@@ -53,6 +53,22 @@ async fn s3_sink_initializes_and_runs(harness: &TestHarness,
fixture: S3SinkFixt
assert_eq!(sinks.len(), 1);
assert_eq!(sinks[0].key, S3_SINK_KEY);
assert!(sinks[0].enabled);
+ assert_eq!(sinks[0].status, ConnectorStatus::Running);
+
+ let keys = fixture.list_objects("").await.expect("List data objects");
+ assert!(
+ keys.is_empty(),
+ "Startup must not publish probe objects: {keys:?}"
+ );
+ let uploads = fixture
+ .bucket()
+ .list_multiparts_uploads(None, None)
+ .await
+ .expect("List multipart uploads");
+ assert!(
+ uploads.iter().all(|page| page.uploads.is_empty()),
+ "Startup must abort its multipart probe: {uploads:?}"
+ );
drop(fixture);
}