This is an automated email from the ASF dual-hosted git repository.

chibenwa pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/james-project.git


The following commit(s) were added to refs/heads/master by this push:
     new f6a00378ec [ENHANCEMENT] Improve IMAP IDLE
f6a00378ec is described below

commit f6a00378ece67c9252b9a14b56c3eb348f052e17
Author: ilya terskov <[email protected]>
AuthorDate: Sat Oct 3 02:59:32 2026 +0700

    [ENHANCEMENT] Improve IMAP IDLE
    
     - Harden IMAP IDLE lifecycle and continuation handling
---
 .../apache/james/imap/api/process/ImapSession.java |  13 +-
 .../james/imap/api/process/SelectedMailbox.java    |  12 +
 .../apache/james/imap/processor/IdleProcessor.java | 221 +++++++++++-----
 .../imap/processor/base/SelectedMailboxImpl.java   |   6 +
 .../imap/processor/IdleProcessorLifecycleTest.java | 290 +++++++++++++++++++++
 .../processor/IdleProcessorSanitizationTest.java   |  62 +++++
 .../processor/base/SelectedMailboxImplTest.java    |  56 ++++
 .../james/imapserver/netty/NettyImapSession.java   |   5 +
 .../james/imapserver/netty/IMAPServerIdleTest.java | 111 ++++++++
 9 files changed, 706 insertions(+), 70 deletions(-)

diff --git 
a/protocols/imap/src/main/java/org/apache/james/imap/api/process/ImapSession.java
 
b/protocols/imap/src/main/java/org/apache/james/imap/api/process/ImapSession.java
index 791e8ce69a..998ae27e4d 100644
--- 
a/protocols/imap/src/main/java/org/apache/james/imap/api/process/ImapSession.java
+++ 
b/protocols/imap/src/main/java/org/apache/james/imap/api/process/ImapSession.java
@@ -127,6 +127,15 @@ public interface ImapSession extends 
CommandDetectionSession {
         return false;
     }
 
+    /**
+     * Return true if the underlying transport connection is currently open 
and active.
+     *
+     * @return true if connected
+     */
+    default boolean isConnected() {
+        return true;
+    }
+
     /**
      * Gets the current client state.
      * 
@@ -237,7 +246,9 @@ public interface ImapSession extends 
CommandDetectionSession {
     boolean startCompression(Runnable runnable);
 
     /**
-     * Push in a new {@link ImapLineHandler} which is called for the next line 
received
+     * Push in a new {@link ImapLineHandler} which is called for the next line 
received.
+     * Implementations must ensure that if an exception is thrown during push,
+     * the handler does not remain installed on the session.
      */
     void pushLineHandler(ImapLineHandler lineHandler);
 
diff --git 
a/protocols/imap/src/main/java/org/apache/james/imap/api/process/SelectedMailbox.java
 
b/protocols/imap/src/main/java/org/apache/james/imap/api/process/SelectedMailbox.java
index bc8b8e7415..54634b2fc3 100644
--- 
a/protocols/imap/src/main/java/org/apache/james/imap/api/process/SelectedMailbox.java
+++ 
b/protocols/imap/src/main/java/org/apache/james/imap/api/process/SelectedMailbox.java
@@ -50,6 +50,18 @@ public interface SelectedMailbox {
 
     void unregisterIdle();
 
+    /**
+     * Unregisters the given IDLE listener only if it matches the currently 
registered listener.
+     * Implementations should override this method to provide identity-safe 
unregistration
+     * (e.g. via atomic compareAndSet) preventing delayed cleanups from 
unregistering newer listeners.
+     * The default implementation falls back to {@link #unregisterIdle()} for 
backward compatibility.
+     *
+     * @param listener the listener instance to unregister
+     */
+    default void unregisterIdle(EventListener.ReactiveEventListener listener) {
+        unregisterIdle();
+    }
+
     boolean isIdling();
 
     /**
diff --git 
a/protocols/imap/src/main/java/org/apache/james/imap/processor/IdleProcessor.java
 
b/protocols/imap/src/main/java/org/apache/james/imap/processor/IdleProcessor.java
index 322d58f2be..b14c1bb763 100644
--- 
a/protocols/imap/src/main/java/org/apache/james/imap/processor/IdleProcessor.java
+++ 
b/protocols/imap/src/main/java/org/apache/james/imap/processor/IdleProcessor.java
@@ -22,10 +22,11 @@ package org.apache.james.imap.processor;
 import static org.apache.james.imap.api.ImapConstants.SUPPORTS_IDLE;
 import static org.apache.james.util.ReactorUtils.logAsMono;
 
+import java.nio.charset.StandardCharsets;
 import java.time.Duration;
 import java.util.List;
 import java.util.Locale;
-import java.util.concurrent.CountDownLatch;
+import java.util.Optional;
 import java.util.concurrent.atomic.AtomicBoolean;
 
 import jakarta.inject.Inject;
@@ -52,16 +53,18 @@ import org.reactivestreams.Publisher;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
-import com.github.fge.lambdas.Throwing;
+import com.google.common.annotations.VisibleForTesting;
 import com.google.common.collect.ImmutableList;
 
 import reactor.core.publisher.Mono;
+import reactor.core.publisher.Sinks;
 
 public class IdleProcessor extends AbstractMailboxProcessor<IdleRequest> 
implements CapabilityImplementingProcessor {
     private static final Logger LOGGER = 
LoggerFactory.getLogger(IdleProcessor.class);
 
     private static final List<Capability> CAPS = 
ImmutableList.of(SUPPORTS_IDLE);
     private static final String DONE = "DONE";
+    private static final int MAX_DISPLAY_LENGTH = 32;
 
     private Duration heartbeatInterval;
     private boolean enableIdle;
@@ -77,87 +80,157 @@ public class IdleProcessor extends 
AbstractMailboxProcessor<IdleRequest> impleme
         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))
-            .then(unsolicitedResponses(session, responder, false))
+        Responder safeResponder = session.threadSafe(responder);
+        SelectedMailbox selectedMailbox = session.getSelected();
+        Sinks.One<Void> idleReadySink = Sinks.one();
+        AtomicBoolean idleActive = new AtomicBoolean(true);
+        AtomicBoolean lineHandlerAdded = new AtomicBoolean(false);
+
+        Optional<IdleMailboxListener> idleListener = 
Optional.ofNullable(selectedMailbox)
+            .map(mailbox -> new IdleMailboxListener(session, safeResponder, 
idleReadySink, idleActive));
+
+        return Mono.fromRunnable(() -> idle(request, session, safeResponder, 
selectedMailbox, idleReadySink, idleActive, lineHandlerAdded, idleListener))
+            .then(unsolicitedResponses(session, safeResponder, false))
             .onErrorResume(e -> {
-                no(request, responder, 
HumanReadableText.GENERIC_FAILURE_DURING_PROCESSING);
+                cleanupIdle(session, selectedMailbox, idleActive, 
lineHandlerAdded, idleReadySink, idleListener);
+                no(request, safeResponder, 
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 boolean cleanupIdle(ImapSession session, SelectedMailbox 
selectedMailbox, AtomicBoolean idleActive,
+                                AtomicBoolean lineHandlerAdded, 
Sinks.One<Void> idleReadySink,
+                                Optional<IdleMailboxListener> idleListener) {
+        boolean cleanupOwner = idleActive.compareAndSet(true, false);
+        if (cleanupOwner) {
+            if (selectedMailbox != null) {
+                idleListener.ifPresent(listener -> {
+                    try {
+                        selectedMailbox.unregisterIdle(listener);
+                    } catch (Exception e) {
+                        LOGGER.debug("Failed to unregister IDLE listener", e);
+                    }
+                });
+            }
+            try {
+                if (session != null && lineHandlerAdded.compareAndSet(true, 
false)) {
+                    session.popLineHandler();
+                }
+            } finally {
+                idleReadySink.tryEmitEmpty();
+            }
+            return true;
         }
+        return false;
+    }
 
-        final AtomicBoolean idleActive = new AtomicBoolean(true);
-
-        session.pushLineHandler((session1, data) -> Mono.fromRunnable(() -> {
-            String line;
-            if (data.length > 2) {
-                line = new String(data, 0, data.length - 2);
+    private void idle(IdleRequest request, ImapSession session, Responder 
safeResponder, SelectedMailbox selectedMailbox,
+                      Sinks.One<Void> idleReadySink, AtomicBoolean idleActive, 
AtomicBoolean lineHandlerAdded,
+                      Optional<IdleMailboxListener> idleListener) {
+        try {
+            if (selectedMailbox != null) {
+                idleListener.ifPresent(selectedMailbox::registerIdle);
             } else {
-                line = "";
+                idleReadySink.tryEmitEmpty();
             }
 
-            if (sm != null) {
-                sm.unregisterIdle();
+            lineHandlerAdded.set(true);
+            try {
+                session.pushLineHandler((session1, data) -> {
+                    if (!idleActive.get()) {
+                        return Mono.empty();
+                    }
+                    cleanupIdle(session1, selectedMailbox, idleActive, 
lineHandlerAdded, idleReadySink, idleListener);
+                    String line = new String(data, 
StandardCharsets.US_ASCII).trim();
+                    if (!session1.isConnected()) {
+                        LOGGER.debug("IDLE continuation received disconnected 
session.");
+                        return Mono.empty();
+                    }
+
+                    if (DONE.equals(line.toUpperCase(Locale.ROOT))) {
+                        okComplete(request, safeResponder);
+                        safeResponder.flush();
+                        return Mono.empty();
+                    }
+
+                    String displayLine = sanitizeForDisplay(line);
+                    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.debug(message);
+                    safeResponder.respond(response);
+                    safeResponder.flush();
+                    return Mono.empty();
+                });
+            } catch (Exception e) {
+                lineHandlerAdded.set(false);
+                throw e;
             }
-            if (!DONE.equals(line.toUpperCase(Locale.US))) {
-                String message = String.format("Continuation for IMAP IDLE was 
not understood. Expected 'DONE', got '%s'.", line);
-                StatusResponse response = getStatusResponseFactory()
-                    .taggedBad(request.getTag(), request.getCommand(),
-                        new 
HumanReadableText("org.apache.james.imap.INVALID_CONTINUATION",
-                            "failed. " + message));
-                LOGGER.info(message);
-                responder.respond(response);
-                responder.flush();
-            } else {
-                okComplete(request, responder);
-                responder.flush();
+
+            // Write continuation response after listener and handler are 
installed (IMAP-341), only if still active
+            if (idleActive.get()) {
+                safeResponder.respond(new 
ContinuationResponse(HumanReadableText.IDLING));
+                safeResponder.flush();
+            }
+
+            if (enableIdle) {
+                scheduleHeartbeat(session, safeResponder, selectedMailbox, 
idleReadySink, idleActive, lineHandlerAdded, idleListener);
             }
-            session1.popLineHandler();
-            idleActive.set(false);
-        }));
-
-        // Check if we should send heartbeats
-        if (enableIdle) {
-            session.schedule(new Runnable() {
-
-                @Override
-                public void run() {
-                    // check if we need to cancel the Runnable
-                    // See IMAP-275
-                    if (session.getState() != ImapSessionState.LOGOUT && 
idleActive.get()) {
-                        // Send a heartbeat to the client to make sure we
-                        // reset the idle timeout. This is kind of the same
-                        // workaround as dovecot use.
-                        //
-                        // This is mostly needed because of the broken
-                        // outlook client, but can't harm for other clients
-                        // too.
-                        // See IMAP-272
+        } catch (Exception e) {
+            cleanupIdle(session, selectedMailbox, idleActive, 
lineHandlerAdded, idleReadySink, idleListener);
+            throw e;
+        }
+    }
+
+    private void scheduleHeartbeat(ImapSession session, Responder 
safeResponder, SelectedMailbox selectedMailbox,
+                                   Sinks.One<Void> idleReadySink, 
AtomicBoolean idleActive, AtomicBoolean lineHandlerAdded,
+                                   Optional<IdleMailboxListener> idleListener) 
{
+        session.schedule(new Runnable() {
+            @Override
+            public void run() {
+                if (session.isConnected() && session.getState() != 
ImapSessionState.LOGOUT && idleActive.get()) {
+                    try {
                         StatusResponse response = 
getStatusResponseFactory().untaggedOk(HumanReadableText.HEARTBEAT);
-                        responder.respond(response);
+                        safeResponder.respond(response);
+                        safeResponder.flush();
 
-                        // schedule the heartbeat again for the next interval
-                        session.schedule(this, heartbeatInterval);
+                        if (idleActive.get() && session.isConnected() && 
session.getState() != ImapSessionState.LOGOUT) {
+                            session.schedule(this, heartbeatInterval);
+                        }
+                    } catch (Exception e) {
+                        LOGGER.debug("Failed to send IMAP IDLE heartbeat, 
stopping keepalive task", e);
+                        cleanupIdle(session, selectedMailbox, idleActive, 
lineHandlerAdded, idleReadySink, idleListener);
                     }
+                } else {
+                    cleanupIdle(session, selectedMailbox, idleActive, 
lineHandlerAdded, idleReadySink, idleListener);
                 }
-            }, heartbeatInterval);
-        }
+            }
+        }, heartbeatInterval);
+    }
 
-        // Write the response after the listener was add
-        // IMAP-341
-        responder.respond(new ContinuationResponse(HumanReadableText.IDLING));
+    @VisibleForTesting
+    static String sanitizeForDisplay(String line) {
+        StringBuilder sanitized = new StringBuilder(Math.min(line.length(), 
MAX_DISPLAY_LENGTH));
+        boolean truncated = false;
+        for (int i = 0; i < line.length(); i++) {
+            if (sanitized.length() == MAX_DISPLAY_LENGTH) {
+                truncated = true;
+                break;
+            }
+            char c = line.charAt(i);
+            if (c >= 32 && c < 127) {
+                sanitized.append(c);
+            }
+        }
+        return truncated ? sanitized + "..." : sanitized.toString();
     }
 
     @Override
@@ -169,12 +242,15 @@ public class IdleProcessor extends 
AbstractMailboxProcessor<IdleRequest> impleme
 
         private final Responder responder;
         private final ImapSession session;
-        private final CountDownLatch countDownLatch;
+        private final Sinks.One<Void> idleReadySink;
+        private final AtomicBoolean idleActive;
 
-        public IdleMailboxListener(ImapSession session, Responder responder, 
CountDownLatch countDownLatch) {
+        public IdleMailboxListener(ImapSession session, Responder responder, 
Sinks.One<Void> idleReadySink,
+                                   AtomicBoolean idleActive) {
             this.session = session;
-            this.responder = session.threadSafe(responder);
-            this.countDownLatch = countDownLatch;
+            this.responder = responder;
+            this.idleReadySink = idleReadySink;
+            this.idleActive = idleActive;
         }
 
         @Override
@@ -184,9 +260,16 @@ public class IdleProcessor extends 
AbstractMailboxProcessor<IdleRequest> impleme
 
         @Override
         public Publisher<Void> reactiveEvent(Event event) {
-            return Mono.fromRunnable(Throwing.runnable(countDownLatch::await))
-                .then(Mono.defer(() -> unsolicitedResponses(session, 
responder, false)))
-                .then(Mono.fromRunnable(responder::flush));
+            return idleReadySink.asMono()
+                .then(Mono.defer(() -> {
+                    if (!idleActive.get()) {
+                        return Mono.empty();
+                    }
+                    return unsolicitedResponses(session, responder, false)
+                        .then(Mono.fromRunnable(responder::flush));
+                }))
+                .onErrorResume(e -> logAsMono(() -> LOGGER.info("Failed to 
push updates to idling client", e)))
+                .then();
         }
 
         @Override
diff --git 
a/protocols/imap/src/main/java/org/apache/james/imap/processor/base/SelectedMailboxImpl.java
 
b/protocols/imap/src/main/java/org/apache/james/imap/processor/base/SelectedMailboxImpl.java
index e922d576ca..7fcff51fe4 100644
--- 
a/protocols/imap/src/main/java/org/apache/james/imap/processor/base/SelectedMailboxImpl.java
+++ 
b/protocols/imap/src/main/java/org/apache/james/imap/processor/base/SelectedMailboxImpl.java
@@ -177,6 +177,11 @@ public class SelectedMailboxImpl implements 
SelectedMailbox, EventListener.React
         idleEventListener.set(null);
     }
 
+    @Override
+    public void unregisterIdle(ReactiveEventListener listener) {
+        idleEventListener.compareAndSet(listener, null);
+    }
+
     @Override
     public boolean isIdling() {
         return idleEventListener.get() != null;
@@ -209,6 +214,7 @@ public class SelectedMailboxImpl implements 
SelectedMailbox, EventListener.React
     }
 
     private synchronized void clearInternalStructures() {
+        idleEventListener.set(null);
         uidMsnConverter.clear();
         flagUpdateUids.clear();
 
diff --git 
a/protocols/imap/src/test/java/org/apache/james/imap/processor/IdleProcessorLifecycleTest.java
 
b/protocols/imap/src/test/java/org/apache/james/imap/processor/IdleProcessorLifecycleTest.java
new file mode 100644
index 0000000000..32842f0b58
--- /dev/null
+++ 
b/protocols/imap/src/test/java/org/apache/james/imap/processor/IdleProcessorLifecycleTest.java
@@ -0,0 +1,290 @@
+/****************************************************************
+ * Licensed to the Apache Software Foundation (ASF) under one   *
+ * or more contributor license agreements.  See the NOTICE file *
+ * distributed with this work for additional information        *
+ * regarding copyright ownership.  The ASF licenses this file   *
+ * to you under the Apache License, Version 2.0 (the            *
+ * "License"); you may not use this file except in compliance   *
+ * with the License.  You may obtain a copy of the License at   *
+ *                                                              *
+ *   http://www.apache.org/licenses/LICENSE-2.0                 *
+ *                                                              *
+ * Unless required by applicable law or agreed to in writing,   *
+ * software distributed under the License is distributed on an  *
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY       *
+ * KIND, either express or implied.  See the License for the    *
+ * specific language governing permissions and limitations      *
+ * under the License.                                           *
+ ****************************************************************/
+
+package org.apache.james.imap.processor;
+
+import static org.apache.james.imap.ImapFixture.TAG;
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+
+import java.nio.charset.StandardCharsets;
+import java.util.ArrayDeque;
+import java.util.Deque;
+import java.util.List;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.function.Consumer;
+
+import org.apache.james.events.EventListener;
+import org.apache.james.imap.api.message.response.ImapResponseMessage;
+import org.apache.james.imap.api.message.response.StatusResponse;
+import org.apache.james.imap.api.process.ImapLineHandler;
+import org.apache.james.imap.api.process.ImapProcessor;
+import org.apache.james.imap.api.process.SelectedMailbox;
+import org.apache.james.imap.encode.FakeImapSession;
+import org.apache.james.imap.message.request.IdleRequest;
+import org.apache.james.imap.message.response.UnpooledStatusResponseFactory;
+import org.apache.james.mailbox.MailboxManager;
+import org.apache.james.metrics.tests.RecordingMetricFactory;
+import org.junit.jupiter.api.Test;
+import org.mockito.ArgumentCaptor;
+
+import reactor.core.publisher.Mono;
+
+class IdleProcessorLifecycleTest {
+
+    private static class TrackingImapSession extends FakeImapSession {
+        final Deque<ImapLineHandler> handlers = new ArrayDeque<>();
+        final AtomicInteger popCount = new AtomicInteger();
+        boolean triggerCallbackDuringPush = false;
+        Consumer<FakeImapSession> onPush;
+
+        @Override
+        public ImapProcessor.Responder threadSafe(ImapProcessor.Responder 
responder) {
+            return responder;
+        }
+
+        @Override
+        public void pushLineHandler(ImapLineHandler lineHandler) {
+            handlers.push(lineHandler);
+            if (triggerCallbackDuringPush) {
+                Mono.from(lineHandler.onLine(this, 
"DONE\r\n".getBytes(StandardCharsets.US_ASCII))).block();
+            }
+            if (onPush != null) {
+                try {
+                    onPush.accept(this);
+                } catch (RuntimeException e) {
+                    handlers.pop();
+                    throw e;
+                }
+            }
+        }
+
+        @Override
+        public void popLineHandler() {
+            popCount.incrementAndGet();
+            if (!handlers.isEmpty()) {
+                handlers.pop();
+            }
+        }
+    }
+
+    private static class RecordingResponder implements ImapProcessor.Responder 
{
+        private final List<ImapResponseMessage> responses = new 
CopyOnWriteArrayList<>();
+
+        @Override
+        public void respond(ImapResponseMessage message) {
+            responses.add(message);
+        }
+
+        @Override
+        public void flush() {
+        }
+
+        public List<ImapResponseMessage> getResponses() {
+            return responses;
+        }
+    }
+
+    @Test
+    void earlyCallbackDuringPushLineHandlerShouldPopHandlerExactlyOnce() {
+        IdleProcessor testee = new IdleProcessor(
+            mock(MailboxManager.class),
+            new UnpooledStatusResponseFactory(),
+            new RecordingMetricFactory());
+
+        TrackingImapSession session = new TrackingImapSession();
+        session.triggerCallbackDuringPush = true;
+
+        ImapLineHandler baseHandler = (session1, data) -> Mono.empty();
+        session.pushLineHandler(baseHandler);
+
+        testee.processRequestReactive(new IdleRequest(TAG), session, new 
RecordingResponder()).block();
+
+        // popLineHandler was invoked exactly once for the IDLE handler, base 
handler preserved
+        assertThat(session.popCount.get()).isEqualTo(1);
+        assertThat(session.handlers).containsExactly(baseHandler);
+    }
+
+    @Test
+    void 
midPushDisconnectOrCleanupShouldPopHandlerExactlyOnceWhenPushCompletes() {
+        IdleProcessor testee = new IdleProcessor(
+            mock(MailboxManager.class),
+            new UnpooledStatusResponseFactory(),
+            new RecordingMetricFactory());
+
+        TrackingImapSession session = new TrackingImapSession();
+        SelectedMailbox selectedMailbox = mock(SelectedMailbox.class);
+        session.selected(selectedMailbox).block();
+
+        ImapLineHandler baseHandler = (session1, data) -> Mono.empty();
+        session.pushLineHandler(baseHandler);
+
+        session.onPush = s -> {
+            ImapLineHandler idleHandler = session.handlers.peek();
+            Mono.from(idleHandler.onLine(s, 
"DONE\r\n".getBytes(StandardCharsets.US_ASCII))).block();
+        };
+
+        testee.processRequestReactive(new IdleRequest(TAG), session, new 
RecordingResponder()).block();
+
+        assertThat(session.popCount.get()).isEqualTo(1);
+        assertThat(session.handlers).containsExactly(baseHandler);
+    }
+
+    @Test
+    void exceptionDuringRegisterIdleShouldAttemptListenerCleanup() {
+        IdleProcessor testee = new IdleProcessor(
+            mock(MailboxManager.class),
+            new UnpooledStatusResponseFactory(),
+            new RecordingMetricFactory());
+
+        TrackingImapSession session = new TrackingImapSession();
+        SelectedMailbox selectedMailbox = mock(SelectedMailbox.class);
+        session.selected(selectedMailbox).block();
+
+        ImapLineHandler baseHandler = (session1, data) -> Mono.empty();
+        session.pushLineHandler(baseHandler);
+
+        doThrow(new RuntimeException("Mailbox error"))
+            .when(selectedMailbox).registerIdle(any());
+
+        RecordingResponder responder = new RecordingResponder();
+        testee.processRequestReactive(new IdleRequest(TAG), session, 
responder).block();
+
+        ArgumentCaptor<EventListener.ReactiveEventListener> captor =
+            ArgumentCaptor.forClass(EventListener.ReactiveEventListener.class);
+        verify(selectedMailbox).unregisterIdle(captor.capture());
+        assertThat(captor.getValue()).isNotNull();
+        assertThat(session.popCount.get()).isZero();
+        assertThat(session.handlers).containsExactly(baseHandler);
+        assertThat(responder.getResponses()).hasSize(1);
+        assertThat(responder.getResponses().get(0))
+            .isInstanceOf(StatusResponse.class);
+        StatusResponse statusResponse = (StatusResponse) 
responder.getResponses().get(0);
+        assertThat(statusResponse.getServerResponseType())
+            .isEqualTo(StatusResponse.Type.NO);
+    }
+
+    @Test
+    void 
exceptionDuringUnregisterIdleShouldStillCleanUpLineHandlerStateAndSink() {
+        IdleProcessor testee = new IdleProcessor(
+            mock(MailboxManager.class),
+            new UnpooledStatusResponseFactory(),
+            new RecordingMetricFactory());
+
+        TrackingImapSession session = new TrackingImapSession();
+        SelectedMailbox selectedMailbox = mock(SelectedMailbox.class);
+        session.selected(selectedMailbox).block();
+
+        doThrow(new RuntimeException("Unregister failed"))
+            .when(selectedMailbox).unregisterIdle(any());
+
+        ImapLineHandler baseHandler = (session1, data) -> Mono.empty();
+        session.pushLineHandler(baseHandler);
+
+        session.triggerCallbackDuringPush = true;
+
+        RecordingResponder responder = new RecordingResponder();
+        testee.processRequestReactive(new IdleRequest(TAG), session, 
responder).block();
+
+        assertThat(session.popCount.get()).isEqualTo(1);
+        assertThat(session.handlers).containsExactly(baseHandler);
+        ArgumentCaptor<EventListener.ReactiveEventListener> captor =
+            ArgumentCaptor.forClass(EventListener.ReactiveEventListener.class);
+        verify(selectedMailbox).unregisterIdle(captor.capture());
+        assertThat(captor.getValue()).isNotNull();
+        assertThat(responder.getResponses()).hasSize(1);
+    }
+
+    @Test
+    void unregisterIdleDuringPushFailureShouldSucceed() {
+        IdleProcessor testee = new IdleProcessor(
+            mock(MailboxManager.class),
+            new UnpooledStatusResponseFactory(),
+            new RecordingMetricFactory());
+
+        TrackingImapSession session = new TrackingImapSession();
+        SelectedMailbox selectedMailbox = mock(SelectedMailbox.class);
+        session.selected(selectedMailbox).block();
+
+        ImapLineHandler baseHandler = (session1, data) -> Mono.empty();
+        session.pushLineHandler(baseHandler);
+
+        session.onPush = s -> {
+            throw new RuntimeException("Push failed");
+        };
+
+        RecordingResponder responder = new RecordingResponder();
+        testee.processRequestReactive(new IdleRequest(TAG), session, 
responder).block();
+
+        ArgumentCaptor<EventListener.ReactiveEventListener> captor =
+            ArgumentCaptor.forClass(EventListener.ReactiveEventListener.class);
+        verify(selectedMailbox).unregisterIdle(captor.capture());
+        assertThat(captor.getValue()).isNotNull();
+        assertThat(session.popCount.get()).isZero();
+        assertThat(session.handlers).containsExactly(baseHandler);
+        assertThat(responder.getResponses()).hasSize(1);
+        assertThat(responder.getResponses().get(0))
+            .isInstanceOf(StatusResponse.class);
+        StatusResponse statusResponse = (StatusResponse) 
responder.getResponses().get(0);
+        assertThat(statusResponse.getServerResponseType())
+            .isEqualTo(StatusResponse.Type.NO);
+    }
+
+    @Test
+    void 
exceptionDuringPopLineHandlerShouldStillCompleteSinkAndPreservePipeline() {
+        IdleProcessor testee = new IdleProcessor(
+            mock(MailboxManager.class),
+            new UnpooledStatusResponseFactory(),
+            new RecordingMetricFactory());
+
+        TrackingImapSession session = new TrackingImapSession() {
+            @Override
+            public void popLineHandler() {
+                super.popLineHandler();
+                throw new RuntimeException("Faulty popLineHandler");
+            }
+        };
+        SelectedMailbox selectedMailbox = mock(SelectedMailbox.class);
+        session.selected(selectedMailbox).block();
+
+        ImapLineHandler baseHandler = (session1, data) -> Mono.empty();
+        session.pushLineHandler(baseHandler);
+
+        session.triggerCallbackDuringPush = true;
+
+        RecordingResponder responder = new RecordingResponder();
+        testee.processRequestReactive(new IdleRequest(TAG), session, 
responder).block();
+
+        assertThat(session.popCount.get()).isEqualTo(1);
+        ArgumentCaptor<EventListener.ReactiveEventListener> captor =
+            ArgumentCaptor.forClass(EventListener.ReactiveEventListener.class);
+        verify(selectedMailbox).unregisterIdle(captor.capture());
+        assertThat(captor.getValue()).isNotNull();
+        assertThat(responder.getResponses()).hasSize(1);
+        assertThat(responder.getResponses().get(0))
+            .isInstanceOf(StatusResponse.class);
+        StatusResponse statusResponse = (StatusResponse) 
responder.getResponses().get(0);
+        assertThat(statusResponse.getServerResponseType())
+            .isEqualTo(StatusResponse.Type.NO);
+    }
+}
diff --git 
a/protocols/imap/src/test/java/org/apache/james/imap/processor/IdleProcessorSanitizationTest.java
 
b/protocols/imap/src/test/java/org/apache/james/imap/processor/IdleProcessorSanitizationTest.java
new file mode 100644
index 0000000000..df9c6d1197
--- /dev/null
+++ 
b/protocols/imap/src/test/java/org/apache/james/imap/processor/IdleProcessorSanitizationTest.java
@@ -0,0 +1,62 @@
+/****************************************************************
+ * Licensed to the Apache Software Foundation (ASF) under one   *
+ * or more contributor license agreements.  See the NOTICE file *
+ * distributed with this work for additional information        *
+ * regarding copyright ownership.  The ASF licenses this file   *
+ * to you under the Apache License, Version 2.0 (the            *
+ * "License"); you may not use this file except in compliance   *
+ * with the License.  You may obtain a copy of the License at   *
+ *                                                              *
+ *   http://www.apache.org/licenses/LICENSE-2.0                 *
+ *                                                              *
+ * Unless required by applicable law or agreed to in writing,   *
+ * software distributed under the License is distributed on an  *
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY       *
+ * KIND, either express or implied.  See the License for the    *
+ * specific language governing permissions and limitations      *
+ * under the License.                                           *
+ ****************************************************************/
+
+package org.apache.james.imap.processor;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+import org.junit.jupiter.api.Test;
+
+class IdleProcessorSanitizationTest {
+
+    @Test
+    void sanitizeForDisplayShouldPreserveNormalCharacters() {
+        assertThat(IdleProcessor.sanitizeForDisplay("DONE"))
+            .isEqualTo("DONE");
+        assertThat(IdleProcessor.sanitizeForDisplay("INVALID_COMMAND"))
+            .isEqualTo("INVALID_COMMAND");
+    }
+
+    @Test
+    void sanitizeForDisplayShouldStripControlCharacters() {
+        assertThat(IdleProcessor.sanitizeForDisplay("LOG\r\nOUT\t\0"))
+            .isEqualTo("LOGOUT");
+    }
+
+    @Test
+    void sanitizeForDisplayShouldTruncateLongInputTo32CharactersWithEllipsis() 
{
+        String longInput = "1234567890123456789012345678901234567890";
+        assertThat(IdleProcessor.sanitizeForDisplay(longInput))
+            .isEqualTo("12345678901234567890123456789012...");
+    }
+
+    @Test
+    void sanitizeForDisplayShouldNotAddEllipsisWhenCleanedStringFitsLimit() {
+        String inputWithManyControlChars = 
"\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0HELLO";
+        assertThat(IdleProcessor.sanitizeForDisplay(inputWithManyControlChars))
+            .isEqualTo("HELLO");
+    }
+
+    @Test
+    @SuppressWarnings("checkstyle:avoidescapedunicodecharacters")
+    void sanitizeForDisplayShouldStripNonAsciiCharacters() {
+        assertThat(IdleProcessor.sanitizeForDisplay("DONE\u200B\u00A0тест"))
+            .isEqualTo("DONE");
+    }
+}
diff --git 
a/protocols/imap/src/test/java/org/apache/james/imap/processor/base/SelectedMailboxImplTest.java
 
b/protocols/imap/src/test/java/org/apache/james/imap/processor/base/SelectedMailboxImplTest.java
index 991177915a..beb69c50d7 100644
--- 
a/protocols/imap/src/test/java/org/apache/james/imap/processor/base/SelectedMailboxImplTest.java
+++ 
b/protocols/imap/src/test/java/org/apache/james/imap/processor/base/SelectedMailboxImplTest.java
@@ -438,4 +438,60 @@ class SelectedMailboxImplTest {
                 .newFlags(newFlags)
                 .build();
     }
+
+    @Test
+    void unregisterIdleWithMatchingListenerShouldClearIdleState() {
+        SelectedMailboxImpl selectedMailbox = new SelectedMailboxImpl(
+            mailboxManager,
+            eventBus,
+            mock(MailboxSession.class),
+            messageManager);
+
+        EventListener.ReactiveEventListener listenerA = 
mock(EventListener.ReactiveEventListener.class);
+        selectedMailbox.registerIdle(listenerA);
+        assertThat(selectedMailbox.isIdling()).isTrue();
+
+        selectedMailbox.unregisterIdle(listenerA);
+        assertThat(selectedMailbox.isIdling()).isFalse();
+    }
+
+    @Test
+    void unregisterIdleWithStaleListenerShouldNotClearNewListener() {
+        SelectedMailboxImpl selectedMailbox = new SelectedMailboxImpl(
+            mailboxManager,
+            eventBus,
+            mock(MailboxSession.class),
+            messageManager);
+
+        EventListener.ReactiveEventListener listenerA = 
mock(EventListener.ReactiveEventListener.class);
+        EventListener.ReactiveEventListener listenerB = 
mock(EventListener.ReactiveEventListener.class);
+
+        // Register initial listener A, then new listener B arrives
+        selectedMailbox.registerIdle(listenerA);
+        selectedMailbox.registerIdle(listenerB);
+
+        // Delayed cleanup of A should NOT unregister B
+        selectedMailbox.unregisterIdle(listenerA);
+        assertThat(selectedMailbox.isIdling()).isTrue();
+
+        // Cleanup of B should successfully unregister
+        selectedMailbox.unregisterIdle(listenerB);
+        assertThat(selectedMailbox.isIdling()).isFalse();
+    }
+
+    @Test
+    void deselectShouldClearIdleListener() {
+        SelectedMailboxImpl selectedMailbox = new SelectedMailboxImpl(
+            mailboxManager,
+            eventBus,
+            mock(MailboxSession.class),
+            messageManager);
+
+        EventListener.ReactiveEventListener listener = 
mock(EventListener.ReactiveEventListener.class);
+        selectedMailbox.registerIdle(listener);
+        assertThat(selectedMailbox.isIdling()).isTrue();
+
+        selectedMailbox.deselect().block();
+        assertThat(selectedMailbox.isIdling()).isFalse();
+    }
 }
diff --git 
a/server/protocols/protocols-imap4/src/main/java/org/apache/james/imapserver/netty/NettyImapSession.java
 
b/server/protocols/protocols-imap4/src/main/java/org/apache/james/imapserver/netty/NettyImapSession.java
index 866ea6d579..9c7acf20f2 100644
--- 
a/server/protocols/protocols-imap4/src/main/java/org/apache/james/imapserver/netty/NettyImapSession.java
+++ 
b/server/protocols/protocols-imap4/src/main/java/org/apache/james/imapserver/netty/NettyImapSession.java
@@ -186,6 +186,11 @@ public class NettyImapSession implements ImapSession, 
NettyConstants {
         return this.selectedMailbox.get();
     }
 
+    @Override
+    public boolean isConnected() {
+        return channel.isActive();
+    }
+
     @Override
     public ImapSessionState getState() {
         return this.state;
diff --git 
a/server/protocols/protocols-imap4/src/test/java/org/apache/james/imapserver/netty/IMAPServerIdleTest.java
 
b/server/protocols/protocols-imap4/src/test/java/org/apache/james/imapserver/netty/IMAPServerIdleTest.java
index d5c5a13206..9874c80850 100644
--- 
a/server/protocols/protocols-imap4/src/test/java/org/apache/james/imapserver/netty/IMAPServerIdleTest.java
+++ 
b/server/protocols/protocols-imap4/src/test/java/org/apache/james/imapserver/netty/IMAPServerIdleTest.java
@@ -27,6 +27,8 @@ import java.nio.ByteBuffer;
 import java.nio.channels.SocketChannel;
 import java.nio.charset.StandardCharsets;
 import java.time.Duration;
+import java.util.List;
+import java.util.regex.Pattern;
 
 import org.apache.james.mailbox.MailboxSession;
 import org.apache.james.mailbox.MessageManager;
@@ -227,4 +229,113 @@ class IMAPServerIdleTest extends AbstractIMAPServerTest {
             assertThat(readStringUntil(clientConnection, s -> s.contains("* 1 
EXISTS")))
                 .isNotNull());
     }
+
+    @Test
+    void invalidContinuationShouldEndIdleAndAllowSubsequentCommands() throws 
Exception {
+        clientConnection.write(ByteBuffer.wrap(String.format("a0 LOGIN %s 
%s\r\n", USER.asString(), USER_PASS).getBytes(StandardCharsets.UTF_8)));
+        readBytes(clientConnection);
+
+        clientConnection.write(ByteBuffer.wrap(("a2 SELECT 
INBOX\r\n").getBytes(StandardCharsets.UTF_8)));
+        readStringUntil(clientConnection, s -> s.contains("a2 OK [READ-WRITE] 
SELECT completed."));
+
+        // Issue IDLE followed by an invalid continuation command
+        clientConnection.write(ByteBuffer.wrap(("a3 
IDLE\r\nINVALID\r\n").getBytes(StandardCharsets.UTF_8)));
+
+        // Expect tagged BAD response for IDLE
+        Awaitility.await().atMost(Duration.ofSeconds(2)).untilAsserted(() ->
+            assertThat(readStringUntil(clientConnection, s -> s.contains("a3 
BAD IDLE failed.")))
+                .isNotNull());
+
+        // Subsequent command must succeed normally, proving line handler was 
cleanly popped
+        clientConnection.write(ByteBuffer.wrap(("a4 
NOOP\r\n").getBytes(StandardCharsets.UTF_8)));
+        Awaitility.await().atMost(Duration.ofSeconds(2)).untilAsserted(() ->
+            assertThat(readStringUntil(clientConnection, s -> s.contains("a4 
OK NOOP completed.")))
+                .isNotNull());
+    }
+
+    @Test
+    void midIdleLogoutShouldRejectContinuationAndAllowSubsequentLogout() 
throws Exception {
+        clientConnection.write(ByteBuffer.wrap(String.format("a0 LOGIN %s 
%s\r\n", USER.asString(), USER_PASS).getBytes(StandardCharsets.UTF_8)));
+        readBytes(clientConnection);
+
+        clientConnection.write(ByteBuffer.wrap(("a2 SELECT 
INBOX\r\n").getBytes(StandardCharsets.UTF_8)));
+        readStringUntil(clientConnection, s -> s.contains("a2 OK [READ-WRITE] 
SELECT completed."));
+
+        clientConnection.write(ByteBuffer.wrap(("a3 
IDLE\r\n").getBytes(StandardCharsets.UTF_8)));
+        readStringUntil(clientConnection, s -> s.contains("+ Idling"));
+
+        // Sending unexpected continuation during IDLE
+        
clientConnection.write(ByteBuffer.wrap(("LOGOUT\r\n").getBytes(StandardCharsets.UTF_8)));
+
+        // Server should reject IDLE with BAD
+        Awaitility.await().atMost(Duration.ofSeconds(2)).untilAsserted(() ->
+            assertThat(readStringUntil(clientConnection, s -> s.contains("a3 
BAD IDLE failed. Continuation for IMAP IDLE was not understood. Expected 
'DONE', got 'LOGOUT'.")))
+                .isNotNull());
+
+        // Subsequent tagged LOGOUT command must succeed normally
+        clientConnection.write(ByteBuffer.wrap(("a4 
LOGOUT\r\n").getBytes(StandardCharsets.UTF_8)));
+        Awaitility.await().atMost(Duration.ofSeconds(2)).untilAsserted(() ->
+            assertThat(readStringUntil(clientConnection, s -> s.contains("a4 
OK LOGOUT completed.")))
+                .isNotNull());
+    }
+
+    @Test
+    void disconnectDuringIdleShouldCleanlyDecrementConnections() throws 
Exception {
+        clientConnection.write(ByteBuffer.wrap(String.format("a0 LOGIN %s 
%s\r\n", USER.asString(), USER_PASS).getBytes(StandardCharsets.UTF_8)));
+        readBytes(clientConnection);
+
+        clientConnection.write(ByteBuffer.wrap(("a2 SELECT 
INBOX\r\n").getBytes(StandardCharsets.UTF_8)));
+        readStringUntil(clientConnection, s -> s.contains("a2 OK [READ-WRITE] 
SELECT completed."));
+
+        clientConnection.write(ByteBuffer.wrap(("a3 
IDLE\r\n").getBytes(StandardCharsets.UTF_8)));
+        readStringUntil(clientConnection, s -> s.contains("+ Idling"));
+
+        // Abruptly sever connection
+        clientConnection.close();
+
+        // Verify connection metric decrements back to 0
+        Awaitility.await().atMost(Duration.ofSeconds(5)).untilAsserted(() ->
+            assertThat(metricFactory.countFor("imapConnections")).isZero());
+    }
+
+    @Test
+    void mailboxEventAfterDoneShouldNotPushUnsolicitedResponses() throws 
Exception {
+        clientConnection.write(ByteBuffer.wrap(String.format("a0 LOGIN %s 
%s\r\n", USER.asString(), USER_PASS).getBytes(StandardCharsets.UTF_8)));
+        readBytes(clientConnection);
+
+        clientConnection.write(ByteBuffer.wrap(("a2 SELECT 
INBOX\r\n").getBytes(StandardCharsets.UTF_8)));
+        readStringUntil(clientConnection, s -> s.contains("a2 OK [READ-WRITE] 
SELECT completed."));
+
+        // Enter IDLE
+        clientConnection.write(ByteBuffer.wrap(("a3 
IDLE\r\n").getBytes(StandardCharsets.UTF_8)));
+        readStringUntil(clientConnection, s -> s.contains("+ Idling"));
+
+        // Complete IDLE cleanly via DONE
+        
clientConnection.write(ByteBuffer.wrap(("DONE\r\n").getBytes(StandardCharsets.UTF_8)));
+        readStringUntil(clientConnection, s -> s.contains("a3 OK IDLE 
completed."));
+
+        // Append message after IDLE is ended
+        inbox.appendMessage(MessageManager.AppendCommand.builder().build("h: 
value\r\n\r\nbody".getBytes()), mailboxSession);
+
+        // Give the now-unregistered IDLE listener a chance to leak an async 
push before
+        // we issue anything else. We deliberately do NOT read the socket 
here: a competing
+        // reader on the same connection could race with and steal bytes from 
the NOOP
+        // response read below, hanging the test. The "exactly once" check on 
EXISTS after
+        // NOOP is what actually detects a leaked duplicate push.
+        Thread.sleep(200);
+
+        // Per RFC 3501 §6.1.2, NOOP legitimately reports pending mailbox 
state changes as
+        // untagged data in its own response - the EXISTS below is expected, 
not a leak.
+        clientConnection.write(ByteBuffer.wrap(("a4 
NOOP\r\n").getBytes(StandardCharsets.UTF_8)));
+        List<String> response = readStringUntil(clientConnection, s -> 
s.contains("a4 OK NOOP completed."));
+        String joinedResponse = String.join("", response);
+        assertThat(joinedResponse).contains("a4 OK NOOP completed.");
+        assertThat(countOccurrences(joinedResponse, "EXISTS"))
+            .describedAs("EXISTS should be reported exactly once by NOOP, not 
duplicated by a stale IDLE push")
+            .isEqualTo(1);
+    }
+
+    private static long countOccurrences(String haystack, String needle) {
+        return 
Pattern.compile(Pattern.quote(needle)).matcher(haystack).results().count();
+    }
 }


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

Reply via email to