chibenwa commented on code in PR #3198:
URL: https://github.com/apache/james-project/pull/3198#discussion_r4103054956
##########
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);
Review Comment:
So we unregister IDLE even if we are not cleanupOwner ?
How does it interact with listenerUnregistered ?
Why do we need two concepts ?
--
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]