Federico Mariani created CAMEL-25167:
----------------------------------------

             Summary: Large payloads - consistent body length detection and 
streaming charset conversion in file / object-storage producers
                 Key: CAMEL-25167
                 URL: https://issues.apache.org/jira/browse/CAMEL-25167
             Project: Camel
          Issue Type: Bug
          Components: camel-ftp, camel-sftp, camel-http, camel-google-storage, 
camel-core, camel-vertx-http, camel-minio, camel-azure, camel-aws2-s3, camel-smb
            Reporter: Federico Mariani
            Assignee: Federico Mariani
             Fix For: 4.23.0


h2. Problem

Several producers handle a large {{InputStream}} / file body inconsistently. 
Depending on the component, the body is copied fully into heap, the upload 
fails, or the exchange never completes. Each component has its own (mostly 
broken) way of finding the body length, and a few convert the whole body to a 
{{String}} just to change the charset.

Confirmed on main (4.23.0-SNAPSHOT):

||Component||Problem||Where||
|aws2-s3, azure-storage-blob|The length is probed with {{mark(1024)}} / 
{{skip(available())}} / {{reset()}}. On a {{BufferedInputStream}} longer than 
the mark buffer (~8 KB) this throws {{IOException: Resetting to invalid mark}}. 
A {{java.nio.file.Path}} body is converted to exactly such a stream 
({{IOConverter.toInputStream(Path)}}), so e.g. the Spring Boot platform-http 
multipart upload (body is a {{Path}}) sent to {{aws2-s3}} fails for anything 
over a few KB. Reproduced with a 2 MB 
file.|{{AWS2S3Utils.determineLengthInputStream}}, 
{{BlobUtils.getInputStreamLength}}|
|azure-storage-blob|A stream whose {{available()}} is 0 (socket-like) is 
measured as length 0.|{{BlobUtils.getInputStreamLength}}|
|azure-storage-datalake|The length is measured by reading the whole stream into 
a {{byte[]}}, then {{reset()}} is called without a 
{{mark()}}.|{{DataLakeUtils}}, {{FileStreamAndLength}}|
|aws2-s3, minio|When the length is unknown, the whole body is copied into a 
{{ByteArrayOutputStream}}, bypassing stream caching (and therefore 
spooling).|{{AWS2S3Producer}} (single put and multipart), {{MinioProducer}}|
|ibm-cos|{{mark(Integer.MAX_VALUE)}} + read to the end + {{reset()}}: the whole 
body is buffered by the stream.|{{IBMCOSProducer}}|
|google-storage|Without the length header, the whole body is read into a 
{{byte[]}} to compute a length that {{Storage.createFrom}} does not need. The 
Content-Length key also ends up in the object's custom 
metadata.|{{GoogleCloudStorageProducer.setContentLength}}|
|vertx-http|A {{WrappedFile}} body that cannot be converted to {{java.io.File}} 
(e.g. sftp/ftp file read into memory, or {{streamDownload=true}} with stream 
caching disabled) sends nothing and never calls the callback: *the exchange 
hangs*. {{Path}} bodies are not handled as files.|{{VertxHttpProducer}} 
process()|
|ftp, sftp, smb|With {{charset}} set, the producer does 
{{getMandatoryBody(String.class).getBytes(charset)}}: the whole file is held in 
heap twice.|{{FtpOperations.storeFile}}, {{SftpOperations}}, {{SmbOperations}}|
|http|{{disableStreamCache=true}} is documented as "use the response stream 
as-is (the stream can only be read once)", but the exchange is not marked, so 
the next processor caches the stream into heap anyway (stream caching is on by 
default, spooling is off).|{{HttpProducer}} / {{DefaultChannel}}|

h2. Proposed fix (minimal, no new options, no default changes)

h3. 1. Shared helpers (camel-util / camel-support)

* *Body length helper*, e.g. {{PayloadHelper.getBodyLength(Exchange)}}: returns 
the length without reading the body, or {{-1}}. Checks, in order: 
{{WrappedFile.getFileLength()}}, {{File}} / {{Path}} size, 
{{StreamCache.length()}}, {{byte[]}} / {{String}} / {{ByteBuffer}}, 
{{ByteArrayInputStream}}, {{FileInputStream}} channel size, then the 
{{CamelFileLength}} header. *No mark/skip/reset probing.*
* *Unknown-length fallback*: when the length is {{-1}} and the SDK needs it, 
copy the stream through {{OutputStreamBuilder.withExchange(exchange)}} (i.e. 
{{CachedOutputStream}}) instead of a {{ByteArrayOutputStream}}. Same behaviour 
as today by default, but it spools to disk when 
{{camel.main.streamCachingSpoolEnabled=true}}, and the length is then read from 
the resulting {{StreamCache}}.
* *Transcoding stream*: generalise {{IOHelper.EncodingInputStream}} (currently 
{{Path}}-only) to any {{InputStream}}, converting from a source charset to a 
target charset chunk by chunk (keeping its surrogate-pair handling).
* *Streaming body helper*, e.g. {{StreamingHelper.setStreamingBody(exchange, 
stream)}}: sets the body and calls 
{{exchange.getExchangeExtension().setStreamCacheDisabled(true)}} (as 
platform-http-vertx and netty-http already do), so a body that is documented as 
read-once is not cached into heap by the next processor.

h3. 2. Component changes

* *aws2-s3, azure-storage-blob, azure-storage-datalake, minio, ibm-cos*: 
replace the per-component length probes and in-heap fallbacks with the two 
helpers above. Treat a {{Path}} body like a {{File}}. Existing length headers 
({{CamelAwsS3ContentLength}}, {{CamelAzureStorageBlobUploadSize}}, ...) keep 
precedence.
* *google-storage*: stop reading the body to compute a length when the length 
is unknown ({{createFrom}} streams without it); stop copying the Content-Length 
key into custom metadata.
* *vertx-http*: when a {{File}} / {{WrappedFile}} / {{Path}} body cannot be 
converted to {{java.io.File}}, fall back to the generic body path instead of 
returning without sending (fixes the hang). Treat {{Path}} like {{File}}.
* *ftp, sftp, smb*: with {{charset}} set, wrap the body stream in the 
transcoding stream (source charset = 
{{ExchangeHelper.getCharsetName(exchange)}}, i.e. what the {{String}} 
conversion uses today, so the written bytes are unchanged). Skip the wrapper 
when source and target charsets are equal.
* *http*: with {{disableStreamCache=true}}, set the response body through the 
streaming body helper, so the behaviour matches the option's documentation.

h3. 3. Tests

* A synthetic, non-markable {{InputStream}} of configurable size (generates 
bytes, stores nothing) in the test support module.
* Per component: a regression test for each problem above (e.g. {{Path}} / 
{{BufferedInputStream}} body > 1 MB to S3 and Azure blob with no length header; 
vertx-http with a {{RemoteFile}} body completes; ftp/sftp/smb {{charset}} 
output byte-for-byte equal to the previous implementation, including multi-byte 
characters across buffer boundaries).
* Where practical, a large-payload test (e.g. several hundred MB) run in a 
forked JVM with a small heap, to show the body is not held in memory.

h2. Out of scope (follow-ups)

* Chunked / multipart upload for unknown lengths in S3 / MinIO / Azure without 
spooling.
* True pass-through for raw platform-http uploads, and streaming multipart 
parsing.
* Streaming responses in vertx-http, knative-http and netty-http.
* Streaming downloads in google-storage, minio and huawei-obs consumers / 
{{getObject}}.
* Making {{streamDownload=true}} (ftp family) and {{includeBody=false}} 
(aws2-s3) disable stream caching: this changes behaviour for routes that rely 
on re-reading the body, so it needs its own discussion and an upgrade-guide 
entry.
* mina-sftp silently ignoring {{charset}}; {{checksumFileAlgorithm}} being 
ignored by remote-file producers; scp buffering the whole file.

h2. Workaround until fixed

Set the length header explicitly (e.g. {{CamelAwsS3ContentLength}} / 
{{CamelAzureStorageBlobUploadSize}} from {{CamelFileLength}}), or convert the 
body to a {{java.io.File}} before the producer. For large bodies, enable 
spooling ({{camel.main.streamCachingSpoolEnabled=true}}) or disable stream 
caching on the route.




--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to