[ 
https://issues.apache.org/jira/browse/CAMEL-25167?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

Federico Mariani updated CAMEL-25167:
-------------------------------------
    Component/s:     (was: camel-vertx-http)
    Description: 
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}}|
|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}}|

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

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

* *Body length helper*, {{PayloadHelper.getBodyLength(Message)}} / 
{{getLength(Object)}}: returns the length without reading the body, or {{-1}}. 
Checks {{WrappedFile.getFileLength()}}, {{File}} / {{Path}} size, 
{{StreamCache.length()}}, {{byte[]}} / {{ByteBuffer}}, {{ByteArrayInputStream}} 
and {{FileInputStream}}. *No mark/skip/reset probing.* The {{CamelFileLength}} 
header is deliberately not used, as it can be stale once the route changed the 
body.
* *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).

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*: determine the length without reading the body, and fall 
back to the stream-caching copy instead of a {{ByteArrayOutputStream}}. The 
Content-Length custom metadata is kept, as existing tests rely on it.
* *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.

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), {{includeBody=false}} (aws2-s3) 
or {{disableStreamCache=true}} (http) disable stream caching for the rest of 
the exchange: 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.

Note on camel-vertx-http: the hang with a remote file body is fixed separately 
by CAMEL-25154.

Note on camel-http {{disableStreamCache=true}}: not changed. Since CAMEL-21162 
the producer returns the raw response stream so that the route's own stream 
caching (with spooling to disk when enabled) handles it; existing tests rely on 
this. To pass the response through without any caching, disable stream caching 
on the route ({{.streamCache("false")}}).

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.


  was:
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}}|

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

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

* *Body length helper*, {{PayloadHelper.getBodyLength(Message)}} / 
{{getLength(Object)}}: returns the length without reading the body, or {{-1}}. 
Checks {{WrappedFile.getFileLength()}}, {{File}} / {{Path}} size, 
{{StreamCache.length()}}, {{byte[]}} / {{ByteBuffer}}, {{ByteArrayInputStream}} 
and {{FileInputStream}}. *No mark/skip/reset probing.* The {{CamelFileLength}} 
header is deliberately not used, as it can be stale once the route changed the 
body.
* *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).

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*: determine the length without reading the body, and fall 
back to the stream-caching copy instead of a {{ByteArrayOutputStream}}. The 
Content-Length custom metadata is kept, as existing tests rely on it.
* *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.

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), {{includeBody=false}} (aws2-s3) 
or {{disableStreamCache=true}} (http) disable stream caching for the rest of 
the exchange: 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.

Note on camel-http {{disableStreamCache=true}}: not changed. Since CAMEL-21162 
the producer returns the raw response stream so that the route's own stream 
caching (with spooling to disk when enabled) handles it; existing tests rely on 
this. To pass the response through without any caching, disable stream caching 
on the route ({{.streamCache("false")}}).

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.



> 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-aws2-s3, camel-azure, camel-core, camel-ftp, 
> camel-google-storage, camel-minio, camel-sftp, camel-smb
>            Reporter: Federico Mariani
>            Assignee: Federico Mariani
>            Priority: Major
>             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}}|
> |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}}|
> h2. Proposed fix (minimal, no new options, no default changes)
> h3. 1. Shared helpers (camel-util / camel-support)
> * *Body length helper*, {{PayloadHelper.getBodyLength(Message)}} / 
> {{getLength(Object)}}: returns the length without reading the body, or 
> {{-1}}. Checks {{WrappedFile.getFileLength()}}, {{File}} / {{Path}} size, 
> {{StreamCache.length()}}, {{byte[]}} / {{ByteBuffer}}, 
> {{ByteArrayInputStream}} and {{FileInputStream}}. *No mark/skip/reset 
> probing.* The {{CamelFileLength}} header is deliberately not used, as it can 
> be stale once the route changed the body.
> * *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).
> 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*: determine the length without reading the body, and fall 
> back to the stream-caching copy instead of a {{ByteArrayOutputStream}}. The 
> Content-Length custom metadata is kept, as existing tests rely on it.
> * *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.
> 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), {{includeBody=false}} 
> (aws2-s3) or {{disableStreamCache=true}} (http) disable stream caching for 
> the rest of the exchange: 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.
> Note on camel-vertx-http: the hang with a remote file body is fixed 
> separately by CAMEL-25154.
> Note on camel-http {{disableStreamCache=true}}: not changed. Since 
> CAMEL-21162 the producer returns the raw response stream so that the route's 
> own stream caching (with spooling to disk when enabled) handles it; existing 
> tests rely on this. To pass the response through without any caching, disable 
> stream caching on the route ({{.streamCache("false")}}).
> 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