prosgarz35 commented on code in PR #3198:
URL: https://github.com/apache/james-project/pull/3198#discussion_r4090699339
##########
protocols/imap/src/main/java/org/apache/james/imap/processor/IdleProcessor.java:
##########
@@ -77,56 +77,83 @@ public void configure(ImapConfiguration imapConfiguration) {
super.configure(imapConfiguration);
this.heartbeatInterval =
imapConfiguration.idleTimeIntervalAsDuration();
- this.enableIdle = imapConfiguration.isEnableIdle();
+ this.enableIdle = imapConfiguration.isEnableIdle() &&
!heartbeatInterval.isZero() && !heartbeatInterval.isNegative();
}
@Override
protected Mono<Void> processRequestReactive(IdleRequest request,
ImapSession session, Responder responder) {
- CountDownLatch countDownLatch = new CountDownLatch(1);
- return Mono.fromRunnable(() -> idle(request, session, responder,
countDownLatch))
+ SelectedMailbox sm = session.getSelected();
+ Sinks.One<Void> idleReadySink = Sinks.one();
+ AtomicBoolean idleActive = new AtomicBoolean(true);
+ AtomicBoolean lineHandlerInstalled = new AtomicBoolean(false);
+ return Mono.fromRunnable(() -> idle(request, session, responder, sm,
idleReadySink, idleActive, lineHandlerInstalled))
.then(unsolicitedResponses(session, responder, false))
.onErrorResume(e -> {
+ cleanupIdle(session, sm, idleActive, lineHandlerInstalled,
idleReadySink);
no(request, responder,
HumanReadableText.GENERIC_FAILURE_DURING_PROCESSING);
return logAsMono(() -> LOGGER.error("Encountered error
executing IMAP IDLE", e));
})
- .then(Mono.fromRunnable(countDownLatch::countDown));
+ .doFinally(signalType -> idleReadySink.tryEmitEmpty());
}
- private void idle(IdleRequest request, ImapSession session, Responder
responder, CountDownLatch countDownLatch) {
- SelectedMailbox sm = session.getSelected();
- if (sm != null) {
- sm.registerIdle(new IdleMailboxListener(session, responder,
countDownLatch));
+ private void cleanupIdle(ImapSession session, SelectedMailbox sm,
AtomicBoolean idleActive,
+ AtomicBoolean lineHandlerInstalled,
Sinks.One<Void> idleReadySink) {
+ if (idleActive.compareAndSet(true, false)) {
+ if (sm != null) {
+ sm.unregisterIdle();
+ }
+ if (lineHandlerInstalled.get()) {
+ session.popLineHandler();
+ }
+ idleReadySink.tryEmitEmpty();
}
+ }
- final AtomicBoolean idleActive = new AtomicBoolean(true);
+ private void idle(IdleRequest request, ImapSession session, Responder
responder, SelectedMailbox sm,
+ Sinks.One<Void> idleReadySink, AtomicBoolean idleActive,
AtomicBoolean lineHandlerInstalled) {
+ if (sm != null) {
+ sm.registerIdle(new IdleMailboxListener(session, sm, responder,
idleReadySink, idleActive, lineHandlerInstalled));
+ } else {
+ idleReadySink.tryEmitEmpty();
+ }
- session.pushLineHandler((session1, data) -> Mono.fromRunnable(() -> {
- String line;
- if (data.length > 2) {
- line = new String(data, 0, data.length - 2);
- } else {
- line = "";
- }
+ try {
+ session.pushLineHandler((session1, data) -> {
+ lineHandlerInstalled.set(true);
+ cleanupIdle(session1, sm, idleActive, lineHandlerInstalled,
idleReadySink);
+ String line = new String(data,
StandardCharsets.US_ASCII).trim();
- if (sm != null) {
- sm.unregisterIdle();
- }
- if (!DONE.equals(line.toUpperCase(Locale.US))) {
- String message = String.format("Continuation for IMAP IDLE was
not understood. Expected 'DONE', got '%s'.", line);
+ if (line.isEmpty() || !session1.isConnected()) {
+ LOGGER.debug("IDLE continuation received empty input or
disconnected session.");
+ return Mono.empty();
+ }
+ String upper = line.toUpperCase(Locale.US);
+ if (DONE.equals(upper)) {
+ okComplete(request, responder);
+ responder.flush();
+ return Mono.empty();
+ }
+ String sanitized = line.replaceAll("[\\r\\n\\x00-\\x1F]", "");
+ String displayLine = sanitized.length() > 32 ?
sanitized.substring(0, 32) + "..." : sanitized;
+ String message = String.format("Continuation for IMAP IDLE was
not understood. Expected 'DONE', got '%s'.", displayLine);
StatusResponse response = getStatusResponseFactory()
.taggedBad(request.getTag(), request.getCommand(),
new
HumanReadableText("org.apache.james.imap.INVALID_CONTINUATION",
"failed. " + message));
- LOGGER.info(message);
+ LOGGER.debug(message);
responder.respond(response);
responder.flush();
- } else {
- okComplete(request, responder);
- responder.flush();
- }
- session1.popLineHandler();
- idleActive.set(false);
- }));
+ return Mono.empty();
+ });
+ lineHandlerInstalled.set(true);
+ } catch (Exception e) {
+ lineHandlerInstalled.set(false);
+ throw e;
+ }
+
+ // Write the response after the listener was added (IMAP-341)
+ responder.respond(new ContinuationResponse(HumanReadableText.IDLING));
+ responder.flush();
Review Comment:
> While I do not disagree with this, can I be provided with an explanation:
>
> * why this responder relocation is needed ?
> * does it impacts netty write ordering in a way not favorable to us?
Hi @chibenwa, here is the rationale and detailed explanation regarding
the placement and `flush()`:
### 1. Why the relocation?
In `upstream/master`, `responder.respond(new ContinuationResponse(...))`
was placed at the very end of `idle()`, after `session.schedule(...)`
(heartbeat registration).
- Relocating it immediately after `session.pushLineHandler(...)`
guarantees that the IDLE lifecycle is initialized in strict logical order:
1. Register the mailbox listener (`registerIdle`).
2. Install the line handler for continuation (`pushLineHandler`).
3. Send `+ Idling` continuation to the client (`respond` + `flush`).
4. Schedule the keepalive heartbeat runnable only once the session is
officially in the IDLE state.
### 2. Does it impact Netty write ordering unfavorably?
**No, in fact it guarantees the correct ordering and prevents buffering
delays:**
1. **Decoder is ready before the wire sees `+ Idling`**:
Both `pushLineHandler(...)` and `responder.respond(...)` execute
synchronously on the connection's Netty EventLoop thread.
Calling `pushLineHandler(...)` first immediately prepends the line
handler to `behaviourOverrides` in `ImapRequestFrameDecoder`. Therefore, the
decoder is fully primed to intercept incoming lines (`DONE\r\n`) *before* the `+
Idling` bytes leave the channel. Even if a fast client sends pipelined
`IDLE\r\nDONE\r\n`, the line handler is already active and will catch `DONE`.
2. **Why `responder.flush()` is needed here**:
`ChannelImapResponseWriter` wraps writes into Netty channel buffers.
In long-running/streaming operations like IDLE, the reactive pipeline does not
complete until the entire command completes (or hangs waiting for input).
Without an immediate `flush()`, the `+ Idling` continuation frame could
sit in the Netty buffer until the pipeline flushes or the buffer fills up,
delaying the client from receiving the continuation cue. Immediate
`responder.flush()`
ensures the client receives `+ Idling\r\n` instantly.
--
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]