prosgarz35 opened a new pull request, #3238:
URL: https://github.com/apache/james-project/pull/3238

   ## Summary of Changes
   
   This PR addresses several critical resource leak vulnerabilities and 
queue-stall scenarios in both `distributed-app` and `postgres-app`:
   
   1. **Reactor `Mono.usingWhen` Resource Leaks on Error and Cancellation**:
      In Project Reactor 3, `Mono.usingWhen(resourceSupplier, resourceClosure, 
asyncCleanup)` only triggers `asyncCleanup` when the publisher completes 
successfully. If the reactive stream terminates with an error (`onError`) or is 
cancelled by a downstream consumer (`onCancel`), the acquired resource is 
**never released** unless the 4-arg or 5-arg `usingWhen` variants provide 
explicit cleanup callbacks.
      This caused unclosed file descriptors, network streams, temporary disk 
files, and R2DBC database connections across multiple core modules:
      - **`S3BlobStoreDAO`**: `InputStream` for S3 upload was left open when an 
S3 upload errored or was cancelled.
      - **`Store` (`blob-common`)**: Inbound HTTP/Blob response `InputStream` 
remained unclosed when reactive byte reading errored or cancelled.
      - **`AESBlobStoreDAO` & `ZstdBlobStoreDAO`**: Temporary 
`FileBackedOutputStream` buffers were not reset/deleted upon 
compression/encryption errors or cancellations.
      - **`MessageStorer`**: Temporary attachment files 
(`parsingResults.dispose()`) were not cleaned up when attachment persistence 
was interrupted or cancelled.
      - **`PostgresTableManager`**: R2DBC connections acquired during extension 
initialization, table migration, listing, truncation, and index initialization 
were not returned to the connection pool if an error occurred or downstream 
cancelled.
   
   2. **RabbitMQ Un-acked Queue Stall on Timeout (`Dequeuer.java`)**:
      In `RabbitMQMailQueue.Dequeuer`, `.timeout(TIMEOUT)` was placed 
downstream of `.onErrorResume(...)`. When blob store fetching stalled and timed 
out, Reactor threw a `TimeoutException` *after* the error handling stage. As a 
result, the RabbitMQ delivery was never rejected/nacked, leaving unacknowledged 
deliveries indefinitely holding worker slots.
      Placing `.timeout(TIMEOUT)` upstream of `.onErrorResume(...)` ensures any 
timeout is caught and properly acknowledged/nacked with consumer backpressure.
   
   ---
   
   ## Detailed Impact & Fixes
   
   ### 1. `server/queue/queue-rabbitmq` - `Dequeuer.java`
   - **Location**: `org.apache.james.queue.rabbitmq.Dequeuer#dequeue`
   - **Problem**: Blob loading timeout produced an unhandled `TimeoutException` 
that bypassed the error-handling and dead-lettering pipeline.
   - **Fix**: Moved `.timeout(TIMEOUT)` before `.onErrorResume(...)` so 
timeouts properly trigger message requeue or nack.
   
   ### 2. `server/blob/blob-s3` - `S3BlobStoreDAO.java`
   - **Location**: `org.apache.james.blob.objectstorage.aws.S3BlobStoreDAO#save`
   - **Problem**: Opening `ByteSource.openStream()` inside `Mono.usingWhen` 
lacked `onError` and `onCancel` rollbacks. Network disruptions during S3 
streaming caused `InputStream` / socket leaks.
   - **Fix**: Added explicit `(stream, error)` and `stream` cleanup functions 
to `Mono.usingWhen`.
   
   ### 3. `server/blob/blob-common` - `Store.java`
   - **Location**: `org.apache.james.blob.api.Store#readByteSource`
   - **Problem**: Reading from `blobStore.readReactive(...)` in `usingWhen` 
leaked response streams on downstream subscriber cancellation or read errors.
   - **Fix**: Added `onError` and `onCancel` stream closure handlers.
   
   ### 4. `server/blob/blob-aes` & `server/blob/blob-zstd`
   - **Location**: `AESBlobStoreDAO#save`, `ZstdBlobStoreDAO#save`, 
`ZstdBlobStoreDAO#compressAndSave`
   - **Problem**: `FileBackedOutputStream` instances were only reset upon 
successful completion, accumulating temporary files in temp storage on failed 
operations.
   - **Fix**: Added `(buffer, error)` and `buffer` cleanup functions to ensure 
`buffer.reset()` is always called on boundedElastic scheduler.
   
   ### 5. `mailbox/store` - `MessageStorer.java`
   - **Location**: 
`org.apache.james.mailbox.store.MessageStorer#storeAttachments`
   - **Problem**: Parsed MIME attachments (`ParsedAttachments`) were not 
disposed if saving metadata to attachment mapper failed or was cancelled.
   - **Fix**: Disposed attachments upon `onError` and `onCancel`.
   
   ### 6. `backends-common/postgres` - `PostgresTableManager.java`
   - **Location**: `initializePostgresExtension`, `initializeTables`, 
`listExistTables`, `truncate`, `initializeTableIndexes`
   - **Problem**: R2DBC connections were not closed if table/index 
initialization steps failed or were cancelled during startup or maintenance 
tasks.
   - **Fix**: Added `onError` and `onCancel` connection close handlers.
   
   ---
   
   ## Verification
   
   All affected modules as well as complete application targets were built and 
verified with Java 21:
   - `backends-common/postgres`: `mvn test-compile` - **SUCCESS**
   - `server/blob/blob-api`, `server/blob/blob-common`, `server/blob/blob-s3`, 
`server/blob/blob-aes`, `server/blob/blob-zstd`: `mvn test-compile` - 
**SUCCESS**
   - `mailbox/store`: `mvn test-compile` - **SUCCESS**
   - `server/queue/queue-rabbitmq`: `mvn test-compile` - **SUCCESS**
   - `server/apps/distributed-app`: `mvn test-compile` - **SUCCESS**
   - `server/apps/postgres-app`: `mvn test-compile` - **SUCCESS**


-- 
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]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to