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]

Reply via email to