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)