Conversation
Codecov Report❌ Patch coverage is 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
🚀 New features to boost your workflow:
|
|
/skill team-review-slim |
There was a problem hiding this comment.
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.
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.
5e385fa to
43d007c
Compare
|
/ready |
| 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(()) |
There was a problem hiding this comment.
State validation checks only offset <= size; an empty key or ETag passes open and can cause recurring poll failures. Reject empty keys and ETags.
| 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(()) |
There was a problem hiding this comment.
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 { |
There was a problem hiding this comment.
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.
| 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))?; |
There was a problem hiding this comment.
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.
| 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; |
There was a problem hiding this comment.
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.
| 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()); |
There was a problem hiding this comment.
The stable-ID prefix is rehashed for every record. Precompute the prefix hash once per batch.
| pub region: Option<String>, | ||
| pub endpoint: Option<String>, | ||
| pub path_style: Option<bool>, | ||
| pub prefix: Option<String>, |
There was a problem hiding this comment.
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)] |
There was a problem hiding this comment.
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 { |
There was a problem hiding this comment.
Add a direct, table-driven ContentRange::from_str test for invalid input branches.
| 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")); |
There was a problem hiding this comment.
The empty-but-truncated listing path and the prefix/StartAfter guard have no protocol coverage. Add wiremock tests for both.
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:
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.