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]