Skip to content

feat(connectors): add S3 source connector - #4347

Open
itsWill wants to merge 1 commit into
apache:masterfrom
itsWill:s3_source_connector
Open

itsWill wants to merge 1 commit into
apache:masterfrom
itsWill:s3_source_connector

Conversation

@itsWill

@itsWill itsWill commented Sep 29, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR address?

Adds the ingestion of an s3 object.

  • reads a bucket from a given path using ListObjectsV2 and StartAfter. Stop at exhaustion instead of continuously rescanning.

  • A restart resumes an scanning the active object, or starts a new walk that can replay completed objects.

  • Read one object at a time with bounded streaming and range requests. Deleted or unreadable active objects block progress rather than being skipped.

  • Treats records as raw bytes separated by an exact configurable delimiter, defaulting to LF. Omit empty records without normalizing payloads.

Closes: #4070

Mimics the bucket walk functionality that redpanda exposes.

What changed?

Ingest delimiter-separated S3 objects:

  • Scan objects from an s3 prefix using a configured delimiter
  • Restarts can replay objects that were scanned
  • Missing or unreadable objects block progress
  • Uses the AWS SDK credential discovery

Local Execution

Passed and tested with a bucket on my aws account.

Also uses the new floci integration tests.

AI Usage

discovery, exploration, and also implementing random bits. Also review.

@codecov

codecov Bot commented Sep 29, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 95.13784% with 97 lines in your changes missing coverage. Please review.
✅ Project coverage is 87.93%. Comparing base (15a49d0) to head (43d007c).
⚠️ Report is 59 commits behind head on master.

Files with missing lines Patch % Lines
core/connectors/sources/s3_source/src/lib.rs 89.55% 26 Missing and 7 partials ⚠️
core/connectors/sources/s3_source/src/client.rs 90.38% 13 Missing and 17 partials ⚠️
core/connectors/sources/s3_source/src/framing.rs 91.41% 26 Missing ⚠️
core/connectors/sources/s3_source/src/reader.rs 97.64% 0 Missing and 6 partials ⚠️
core/connectors/sources/s3_source/src/state.rs 99.21% 0 Missing and 2 partials ⚠️
Additional details and impacted files
@@             Coverage Diff              @@
##             master    #4347      +/-   ##
============================================
+ Coverage     87.56%   87.93%   +0.37%     
- Complexity     1575     1578       +3     
============================================
  Files          1284     1298      +14     
  Lines        225457   230839    +5382     
  Branches     188821   194189    +5368     
============================================
+ Hits         197413   202992    +5579     
+ Misses        23309    23073     -236     
- Partials       4735     4774      +39     
Components Coverage Δ
Rust Core 89.03% <95.13%> (+0.27%) ⬆️
Java SDK 68.74% <ø> (+0.06%) ⬆️
C# SDK 77.53% <ø> (+0.10%) ⬆️
Python SDK 91.24% <ø> (+0.26%) ⬆️
PHP SDK 85.67% <ø> (ø)
Node SDK 96.59% <ø> (+1.84%) ⬆️
Go SDK 70.32% <ø> (+0.27%) ⬆️
Files with missing lines Coverage Δ
core/connectors/sources/s3_source/src/batch.rs 100.00% <100.00%> (ø)
core/connectors/sources/s3_source/src/config.rs 100.00% <100.00%> (ø)
core/connectors/sources/s3_source/src/message.rs 100.00% <100.00%> (ø)
core/connectors/sources/s3_source/src/state.rs 99.21% <99.21%> (ø)
core/connectors/sources/s3_source/src/reader.rs 97.64% <97.64%> (ø)
core/connectors/sources/s3_source/src/framing.rs 91.41% <91.41%> (ø)
core/connectors/sources/s3_source/src/client.rs 90.38% <90.38%> (ø)
core/connectors/sources/s3_source/src/lib.rs 89.55% <89.55%> (ø)

... and 207 files with indirect coverage changes

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.
  • 📦 JS Bundle Analysis: Save yourself from yourself by tracking and limiting bundle sizes in JS merges.

@itsWill
itsWill marked this pull request as ready for review September 30, 2026 11:46
@github-actions github-actions Bot added the S-waiting-on-review PR is waiting on a reviewer label Sep 30, 2026
@hubcio

hubcio commented Oct 1, 2026

Copy link
Copy Markdown
Contributor

/skill team-review-slim

@github-actions github-actions Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Summary: The review of the new S3 source connector found two operational problems. A missing or unreadable restored active object fails open() and leaves the connector skipped until an operator restarts it, and the new crate is absent from the version-bump script.

Counts: critical 0, warning 2, nit 3, simplification 2


This review was generated by Claude Code 2.1.284 on deepseek-flash[1m]. Review the output before you act on it.

Comment thread core/connectors/sources/s3_source/src/lib.rs Outdated
Comment thread core/connectors/sources/s3_source/Cargo.toml
Comment thread .github/workflows/_build_rust_artifacts.yml
Comment thread core/connectors/sources/s3_source/src/client.rs
Comment thread core/connectors/sources/s3_source/src/state.rs Outdated
Comment thread core/connectors/sources/s3_source/src/client.rs Outdated
Comment thread core/connectors/sources/s3_source/src/state.rs
@github-actions github-actions Bot added S-waiting-on-author PR is waiting on author response and removed S-waiting-on-review PR is waiting on a reviewer labels Oct 1, 2026
Support ingestion of delimiter-separated S3 objects without requiring
producers to publish directly to Iggy or modifying their source objects.

Keep discovery to a finite lexicographical prefix walk using ListObjectsV2
and StartAfter. Stop at exhaustion instead of continuously rescanning.
Persist only the active object key, ETag, size and acknowledged byte
offset, avoiding a growing completed-key set. The cursor and exhaustion
remain transient: restart resumes an active object, or starts a new walk
that can replay completed objects.

Read one object at a time with bounded streaming and range requests.
Use conditional reads to prevent applying saved offsets to changed
contents, restarting an active replacement at byte zero. Deleted or
unreadable active objects block progress rather than being skipped.

Treat records as raw bytes separated by an exact configurable delimiter,
defaulting to LF. Omit empty records without normalizing payloads. Leave
compression, format parsing, suffix filters, parallel object reads and
notification-driven ingestion outside this scope. Never tag, move or
delete source objects. Use AWS SDK credential discovery rather than
plugin-specific access-key fields.
@itsWill
itsWill force-pushed the s3_source_connector branch from 5e385fa to 43d007c Compare October 2, 2026 11:36
@itsWill

itsWill commented Oct 2, 2026

Copy link
Copy Markdown
Contributor Author

/ready

@github-actions github-actions Bot added S-waiting-on-review PR is waiting on a reviewer and removed S-waiting-on-author PR is waiting on author response labels Oct 2, 2026
Comment on lines +139 to +147
pub(crate) fn validate(&self) -> Result<(), StateError> {
if self.next_byte_offset > self.size {
return Err(StateError::OffsetBeyondObjectSize {
next_byte_offset: self.next_byte_offset,
size: self.size,
});
}

Ok(())

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

State validation checks only offset <= size; an empty key or ETag passes open and can cause recurring poll failures. Reject empty keys and ETags.

Comment on lines +157 to +163
pub(crate) async fn validate_access(client: &Client, config: &ResolvedConfig) -> Result<(), Error> {
if let Some(mut object) = list_next_object(client, config, None).await? {
// GET, rather than HEAD alone, also checks permission to decrypt/read.
// Drop the body: startup never consumes or checkpoints records.
open_object(client, config, &mut object).await?;
}
Ok(())

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The access probe opens a GET for the first object and drops the streaming body unread. This can waste transferred bytes and prevent connection reuse. Use a Range: bytes=0-0 GET with an empty-object guard so the probe still verifies read and decrypt permission without opening an unbounded response body.

state::{SourceState, StateTracker},
};

async fn source(server: &MockServer, checkpoint: Option<ConnectorState>) -> S3Source {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The source helper bypasses open, leaving the fresh-start access probe without wiremock coverage. Exercise fresh open for successful probing, an empty listing, and denied access.

Comment on lines +58 to +67
let response = client
.list_objects_v2()
.bucket(&config.bucket)
.set_prefix(config.prefix.clone())
.set_start_after(start_after.map(str::to_owned))
.encoding_type(EncodingType::Url)
.max_keys(1)
.send()
.await
.map_err(|error| request_error("list", &error))?;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

max_keys(1) requires one LIST round trip per object and makes scans of small objects latency-bound. Fetch larger pages and maintain a local queue, or document the throughput ceiling.

Comment on lines +142 to +147
object.etag = metadata
.e_tag()
.ok_or_else(|| invalid_response("missing ETag"))?
.to_owned();
object.size = object_size(metadata.content_length())?;
object.next_byte_offset = 0;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

A 412 refresh changes only the local object and then loses the refreshed identity on error, so every later poll repeats GET(412) plus HEAD. Cache the refreshed identity across polls and retain the pre-call offset for error context.

Comment on lines +29 to +36
let mut hasher = blake3::Hasher::new();

for part in [endpoint, bucket, key, etag] {
hasher.update(&(part.len() as u64).to_le_bytes());
hasher.update(part.as_bytes());
}

hasher.update(&record_start_offset.to_le_bytes());

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The stable-ID prefix is rehashed for every record. Precompute the prefix hash once per batch.

Comment on lines +26 to +29
pub region: Option<String>,
pub endpoint: Option<String>,
pub path_style: Option<bool>,
pub prefix: Option<String>,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Empty region and endpoint strings are accepted and fail later with obscure errors; an empty prefix is retained even though it is semantically equivalent to no prefix. Reject an empty region or endpoint and normalize an empty prefix to None.

use iggy_common::{IGGY_MESSAGE_HEADER_SIZE, MAX_MESSAGE_SIZE_UPPER_BYTES, MAX_PAYLOAD_SIZE};
use serde::Deserialize;

#[derive(Debug, Deserialize)]

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Unknown fields are silently ignored, so configuration typos can pass unnoticed. Check compatibility with sibling connectors, then add deny_unknown_fields or supported aliases.

total: u64,
}

impl FromStr for ContentRange {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Add a direct, table-driven ContentRange::from_str test for invalid input branches.

Comment on lines +69 to +95
let Some(object) = response.contents().first() else {
// Without a grouping delimiter an empty, truncated page cannot move
// this key-based cursor forward. Do not mistake it for exhaustion.
if response.is_truncated() == Some(true) {
return Err(Error::HttpRequestFailed(
"S3 returned an empty truncated listing".into(),
));
}
return Ok(None);
};
let key = object
.key()
.ok_or_else(|| invalid_response("missing key"))?;
let key = if response.encoding_type() == Some(&EncodingType::Url) {
percent_decode_str(key)
.decode_utf8()
.map_err(|_| invalid_response("key is not valid UTF-8"))?
} else {
Cow::Borrowed(key)
};
if start_after.is_some_and(|cursor| key.as_ref() <= cursor)
|| config
.prefix
.as_ref()
.is_some_and(|prefix| !key.starts_with(prefix))
{
return Err(invalid_response("listing does not honor prefix/StartAfter"));

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The empty-but-truncated listing path and the prefix/StartAfter guard have no protocol coverage. Add wiremock tests for both.

@github-actions github-actions Bot added S-waiting-on-author PR is waiting on author response and removed S-waiting-on-review PR is waiting on a reviewer labels Oct 5, 2026

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

S-waiting-on-author PR is waiting on author response

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Implement the source connector for s3

3 participants