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

Work on CAMEL-25167 started by Federico Mariani.
------------------------------------------------
> 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-http, camel-minio, camel-sftp, camel-smb, 
> camel-vertx-http
>            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}}|
> |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