quantranhong1999 commented on code in PR #3198:
URL: https://github.com/apache/james-project/pull/3198#discussion_r4103441124


##########
protocols/imap/src/main/java/org/apache/james/imap/processor/IdleProcessor.java:
##########
@@ -77,87 +88,264 @@ 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))
-            .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);
+        AtomicReference<LineHandlerState> lineHandlerState = new 
AtomicReference<>(LineHandlerState.NOT_INSTALLED);
+        AtomicReference<EventListener.ReactiveEventListener> idleListenerRef = 
new AtomicReference<>();
+        AtomicBoolean listenerUnregistered = new AtomicBoolean(false);
+        return Mono.fromRunnable(() -> idle(request, session, safeResponder, 
selectedMailbox, idleReadySink, idleActive, lineHandlerState, idleListenerRef, 
listenerUnregistered))
+            .then(unsolicitedResponses(session, safeResponder, false))
             .onErrorResume(e -> {
-                no(request, responder, 
HumanReadableText.GENERIC_FAILURE_DURING_PROCESSING);
+                cleanupIdle(session, selectedMailbox, idleActive, 
lineHandlerState, idleReadySink, idleListenerRef.get(), listenerUnregistered);
+                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,
+                                AtomicReference<LineHandlerState> 
lineHandlerState, Sinks.One<Void> idleReadySink,
+                                EventListener.ReactiveEventListener 
idleListener, AtomicBoolean listenerUnregistered) {
+        boolean cleanupOwner = idleActive.compareAndSet(true, false);
+        try {
+            unregisterIdleOnce(selectedMailbox, idleListener, 
listenerUnregistered);
+        } catch (Exception e) {
+            LOGGER.debug("Failed to unregister IDLE listener", e);
+        }
+        if (cleanupOwner) {
+            LineHandlerState previous = lineHandlerState.getAndUpdate(state -> 
{
+                if (state == LineHandlerState.INSTALLING) {
+                    return LineHandlerState.REMOVAL_PENDING;
+                }
+                if (state == LineHandlerState.INSTALLED) {
+                    return LineHandlerState.REMOVED;
+                }
+                return state;
+            });
+            try {
+                if (previous == LineHandlerState.INSTALLED) {
+                    session.popLineHandler();
+                }
+            } finally {
+                idleReadySink.tryEmitEmpty();
+            }
+            return true;
         }
+        return false;
+    }
 
-        final AtomicBoolean idleActive = new AtomicBoolean(true);
+    private void unregisterIdleOnce(SelectedMailbox selectedMailbox, 
EventListener.ReactiveEventListener idleListener,
+                                    AtomicBoolean listenerUnregistered) {
+        if (selectedMailbox == null || idleListener == null || 
listenerUnregistered == null) {
+            return;
+        }
+        if (listenerUnregistered.compareAndSet(false, true)) {
+            try {
+                selectedMailbox.unregisterIdle(idleListener);
+            } catch (Exception e) {
+                listenerUnregistered.set(false);
+                throw e;
+            }
+        }
+    }
 
-        session.pushLineHandler((session1, data) -> Mono.fromRunnable(() -> {
-            String line;
-            if (data.length > 2) {
-                line = new String(data, 0, data.length - 2);
+    private void cleanupUnregisteredIdle(SelectedMailbox selectedMailbox, 
EventListener.ReactiveEventListener idleListener,
+                                         AtomicReference<LineHandlerState> 
lineHandlerState, Sinks.One<Void> idleReadySink,
+                                         AtomicBoolean listenerUnregistered) {
+        try {
+            unregisterIdleOnce(selectedMailbox, idleListener, 
listenerUnregistered);
+        } finally {
+            lineHandlerState.compareAndSet(LineHandlerState.REMOVAL_PENDING, 
LineHandlerState.REMOVED);
+            idleReadySink.tryEmitEmpty();
+        }
+    }
+
+    private EventListener.ReactiveEventListener 
registerIdleListener(ImapSession session, Responder safeResponder,
+                                                                    
SelectedMailbox selectedMailbox, Sinks.One<Void> idleReadySink,
+                                                                    
AtomicBoolean idleActive, AtomicReference<LineHandlerState> lineHandlerState,
+                                                                    
AtomicReference<EventListener.ReactiveEventListener> idleListenerRef,
+                                                                    
AtomicBoolean listenerUnregistered) {
+        EventListener.ReactiveEventListener idleListener = null;
+        try {
+            if (selectedMailbox != null) {
+                idleListener = new IdleMailboxListener(session, 
selectedMailbox, safeResponder, idleReadySink, idleActive, lineHandlerState, 
listenerUnregistered);
+                idleListenerRef.set(idleListener);
+                selectedMailbox.registerIdle(idleListener);
             } else {
-                line = "";
+                idleReadySink.tryEmitEmpty();
             }
 
-            if (sm != null) {
-                sm.unregisterIdle();
+            if (!idleActive.get()) {
+                cleanupUnregisteredIdle(selectedMailbox, idleListener, 
lineHandlerState, idleReadySink, listenerUnregistered);
+                return null;
             }
-            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();
+            return idleListener;
+        } catch (Exception e) {
+            try {
+                unregisterIdleOnce(selectedMailbox, idleListener, 
listenerUnregistered);
+            } catch (Exception cleanupException) {
+                e.addSuppressed(cleanupException);
+            } finally {
+                lineHandlerState.set(LineHandlerState.REMOVED);
+            }
+            throw e;
+        }
+    }
+
+    private ImapLineHandler createIdleLineHandler(IdleRequest request, 
Responder safeResponder, SelectedMailbox selectedMailbox,
+                                                  Sinks.One<Void> 
idleReadySink, AtomicBoolean idleActive,
+                                                  
AtomicReference<LineHandlerState> lineHandlerState,
+                                                  
EventListener.ReactiveEventListener idleListener,
+                                                  AtomicBoolean 
listenerUnregistered) {
+        return (session1, data) -> {
+            if (!idleActive.get()) {
+                return Mono.empty();
+            }
+            lineHandlerState.compareAndSet(LineHandlerState.INSTALLING, 
LineHandlerState.INSTALLED);
+            if (!cleanupIdle(session1, selectedMailbox, idleActive, 
lineHandlerState, idleReadySink, idleListener, listenerUnregistered)) {
+                // IDLE was already cleaned up by another thread (heartbeat, 
disconnect, etc.)
+                return Mono.empty();
+            }
+            String line = new String(data, StandardCharsets.US_ASCII).trim();
+
+            if (line.isEmpty() || !session1.isConnected()) {
+                LOGGER.debug("IDLE continuation received empty input or 
disconnected session.");
+                return Mono.empty();
+            }

Review Comment:
   `cleanupIdle` has already run at this point: the line handler is removed and 
`idleActive` is false. So if the client sends an empty line, we return without 
sending any response for the IDLE tag. The client never gets a reply to `a3 
IDLE`. Before this PR it got `BAD ... got ''`. Could we return early only when 
the session is disconnected, and still answer the empty line with BAD?



##########
protocols/imap/src/main/java/org/apache/james/imap/processor/IdleProcessor.java:
##########
@@ -184,9 +382,19 @@ public boolean isHandling(Event event) {
 
         @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 -> {
+                    cleanupIdle(session, selectedMailbox, idleActive, 
lineHandlerState, idleReadySink, this, listenerUnregistered);

Review Comment:
   I don't think a failure while pushing one event should end the whole IDLE. 
This removes the line handler and unregisters the listener, but sends no tagged 
response. The client thinks it's still idling and gets no more updates. When it 
finally sends `DONE`, that's parsed as a new command and gets a BAD. Before, an 
error here only lost that one notification. Could we just log and keep IDLE 
running?



##########
protocols/imap/src/main/java/org/apache/james/imap/processor/base/SelectedMailboxImpl.java:
##########
@@ -402,7 +408,7 @@ public void resetNewApplicableFlags() {
     public Publisher<Void> reactiveEvent(Event event) {
         return Mono.fromRunnable(() -> synchronizedEvent(event))
             .subscribeOn(Schedulers.boundedElastic())
-            .then(Mono.fromCallable(idleEventListener::get)
+            .then(Mono.defer(() -> Mono.justOrEmpty(idleEventListener.get()))

Review Comment:
   `Mono.fromCallable` already gives an empty Mono when the value is `null`. I 
do not think this change is needed.



-- 
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