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);
 }

Reply via email to