leaves12138 commented on code in PR #989: URL: https://github.com/apache/paimon-rust/pull/989#discussion_r4124009385
########## crates/paimon/src/io/uri_reader.rs: ########## @@ -0,0 +1,353 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! URI dispatch for Blob references, matching Java's UriReaderFactory. +//! Ordinary paths retain the table's FileIO, including provider credentials. + +use super::{FileIO, FileRead, InputFile}; +use crate::{Error, Result}; +use bytes::{Buf, Bytes, BytesMut}; +use reqwest::header::CONTENT_ENCODING; +use std::ops::Range; +use std::sync::{Arc, LazyLock}; + +pub(crate) enum UriInput { + File(InputFile), + Http(HttpReader), +} + +impl UriInput { + pub(crate) fn new(file_io: &FileIO, uri: &str) -> Result<Self> { + let http = uri.split_once(':').is_some_and(|(scheme, _)| { + scheme.eq_ignore_ascii_case("http") || scheme.eq_ignore_ascii_case("https") + }); + if http { + Ok(Self::Http(HttpReader { uri: uri.into() })) + } else { + Ok(Self::File(file_io.new_input(uri)?)) + } + } + + pub(crate) async fn size(&self) -> Result<u64> { + match self { + Self::File(input) => Ok(input.metadata().await?.size), + Self::Http(reader) => reader.size().await, + } + } + + pub(crate) async fn reader(&self) -> Result<Arc<dyn FileRead>> { + match self { + Self::File(input) => Ok(Arc::new(input.reader().await?)), + Self::Http(reader) => Ok(Arc::new(reader.clone())), + } + } + + /// Reuse one HTTP response while copying a descriptor into a Blob file. + /// Reopening a Range-ignoring endpoint per chunk would reread every prefix. + pub(crate) async fn reader_for_range(&self, range: Range<u64>) -> Result<Arc<dyn FileRead>> { + match self { + Self::File(input) => Ok(Arc::new(input.reader().await?)), + Self::Http(reader) => Ok(Arc::new(reader.open(range).await?)), + } + } +} + +static HTTP_CLIENT: LazyLock<reqwest::Client> = LazyLock::new(reqwest::Client::new); + +#[derive(Clone)] +pub(crate) struct HttpReader { + uri: String, +} + +impl HttpReader { + fn request(&self) -> reqwest::RequestBuilder { + HTTP_CLIENT.get(&self.uri) + } + + async fn response(&self) -> Result<reqwest::Response> { + let response = self + .request() + .send() + .await + .map_err(http_error)? + .error_for_status() + .map_err(http_error)?; + if response.status() != reqwest::StatusCode::OK { + return Err(invalid(&format!( + "Unexpected HTTP Blob status: {}", + response.status() + ))); + } + // Reqwest removes Content-Encoding after decoding gzip/deflate. + // Never silently copy an unsupported encoded representation as payload. + if response + .headers() + .get(CONTENT_ENCODING) + .is_some_and(|value| value.as_bytes() != b"identity") + { + return Err(invalid("Unsupported HTTP Blob Content-Encoding")); + } + Ok(response) + } + + async fn size(&self) -> Result<u64> { + // GET also works with endpoints that reject HEAD. A chunked or decoded + // response has no known size: count without retaining the payload. + let mut response = self.response().await?; + if let Some(size) = response.content_length() { + return Ok(size); + } + let mut size = 0_u64; + while let Some(chunk) = response.chunk().await.map_err(http_error)? { + size = size + .checked_add(chunk.len() as u64) + .ok_or_else(|| invalid("HTTP Blob size exceeds u64"))?; + } + Ok(size) + } +} + +impl HttpReader { + async fn open(&self, range: Range<u64>) -> Result<HttpRangeReader> { + if range.start >= range.end { + return Err(invalid("Invalid HTTP Blob range")); + } + // Like Java's HttpUriReader, offsets address the decoded entity. + // A wire Range may instead address compressed bytes, so use a decoded + // GET stream and skip the prefix once for the entire copy. + let response = self.response().await?; + let skip = range.start; + Ok(HttpRangeReader { + end: range.end, + state: tokio::sync::Mutex::new(HttpRangeState { + response, + position: range.start, + skip, + pending: Bytes::new(), + }), + }) + } +} + +#[async_trait::async_trait] +impl FileRead for HttpReader { + async fn read(&self, range: Range<u64>) -> Result<Bytes> { + if range.start == range.end { + return Ok(Bytes::new()); + } + self.open(range.clone()).await?.read(range).await + } +} + +struct HttpRangeReader { + end: u64, + state: tokio::sync::Mutex<HttpRangeState>, +} + +struct HttpRangeState { + response: reqwest::Response, + position: u64, + skip: u64, + pending: Bytes, +} + +#[async_trait::async_trait] +impl FileRead for HttpRangeReader { + async fn read(&self, range: Range<u64>) -> Result<Bytes> { + let mut state = self.state.lock().await; + if range.start != state.position || range.end < range.start || range.end > self.end { + return Err(invalid("HTTP Blob copy requires consecutive ranges")); + } + let length = range.end - range.start; + let mut remaining = length; + let mut result = BytesMut::new(); + while remaining > 0 { + if state.pending.is_empty() { + let Some(chunk) = state.response.chunk().await.map_err(http_error)? else { + return Err(invalid(&format!( + "HTTP Blob short read: expected {length} bytes, received {}", + length - remaining + ))); + }; + state.pending = chunk; + } + let skip = state.skip.min(state.pending.len() as u64) as usize; + state.skip -= skip as u64; + state.pending.advance(skip); + let count = remaining.min(state.pending.len() as u64) as usize; + result.extend_from_slice(&state.pending[..count]); + state.pending.advance(count); + remaining -= count as u64; + } + state.position = range.end; + Ok(result.freeze()) + } +} + +fn http_error(error: reqwest::Error) -> Error { + Error::UnexpectedError { + message: format!("HTTP Blob request failed: {error}"), Review Comment: [P1] Redact signed HTTP URLs from both the error message and cause `reqwest::Error` includes the request/response URL in its display output, and retaining it as `source` also exposes that URL when the error chain is formatted. I reproduced this through the public `BlobReader::from_file_io(...).read_blobs(...)` API: an unsigned `/redirect` descriptor redirects to `/failure?signature=review-synthetic-secret`, which returns HTTP 403. The returned Paimon error contains `signature=review-synthetic-secret` in both the message and the nested reqwest error. `blob_error_with_context` does not protect this case because it only replaces the original descriptor URI, not a different redirect target. A signed URL is a credential, so ordinary HTTP failures would expose it in query/job logs. Please ensure HTTP errors and their causes never retain unredacted userinfo/query credentials, including redirected URLs, and sanitize the URI context added by the table read/write wrappers as well. Java's `HttpClientUtils` explicitly sanitizes HTTP failures rather than propagating the original exception. A redirect-to-signed-URL failure test should cover this path. ########## crates/paimon/src/table/table_write.rs: ########## @@ -620,6 +618,12 @@ impl TableWrite { if batch.num_rows() == 0 { return Ok(None); } + if !self.primary_key_indices.is_empty() { + super::inline_blob::validate_inline_blob_columns( Review Comment: [P2] Apply row-kind filtering before validating inline Blob payloads The new validation runs before `enrich_rowkind_batch`, which is where `filter_rowkind_batch` discards ignored rows. With a primary-key table configured with `blob-descriptor-field=payload` and `ignore-delete=true`, a batch containing an INSERT with a valid descriptor and a DELETE with an empty payload now fails with `BlobDescriptor bytes too short`, instead of dropping the DELETE and writing the INSERT. I reproduced this with explicit `_VALUE_KIND` values `[0, 3]`; the same integration test passes on parent `fc27163` and fails on this head. Ignored UPDATE_BEFORE rows are subject to the same ordering issue. Java's `TableWriteImpl.writeAndReturn` returns for a filtered-out row before extracting or writing its Blob value. Please validate only the rows that survive row-kind generation/filtering, still before opening any physical files, and add a mixed retained/ignored-row regression test. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
