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]