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]