chibenwa commented on code in PR #3198:
URL: https://github.com/apache/james-project/pull/3198#discussion_r4116951487
##########
protocols/imap/src/main/java/org/apache/james/imap/processor/IdleProcessor.java:
##########
@@ -77,87 +88,235 @@ 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<>();
+ return Mono.fromRunnable(() -> idle(request, session, safeResponder,
selectedMailbox, idleReadySink, idleActive, lineHandlerState, idleListenerRef))
+ .then(unsolicitedResponses(session, safeResponder, false))
.onErrorResume(e -> {
- no(request, responder,
HumanReadableText.GENERIC_FAILURE_DURING_PROCESSING);
+ cleanupIdle(session, selectedMailbox, idleActive,
lineHandlerState, idleReadySink, idleListenerRef.get());
+ 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) {
+ boolean cleanupOwner = idleActive.compareAndSet(true, false);
+ if (cleanupOwner) {
+ if (selectedMailbox != null && idleListener != null) {
+ try {
+ selectedMailbox.unregisterIdle(idleListener);
+ } catch (Exception e) {
+ LOGGER.debug("Failed to unregister IDLE listener", e);
+ }
+ }
+ 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 != null)
{
+ 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 EventListener.ReactiveEventListener
registerIdleListener(ImapSession session, Responder safeResponder,
+
SelectedMailbox selectedMailbox, Sinks.One<Void> idleReadySink,
+
AtomicBoolean idleActive, AtomicReference<LineHandlerState> lineHandlerState,
+
AtomicReference<EventListener.ReactiveEventListener> idleListenerRef) {
+ EventListener.ReactiveEventListener idleListener = null;
+ try {
+ if (selectedMailbox != null) {
+ idleListener = new IdleMailboxListener(session, safeResponder,
idleReadySink, idleActive);
+ idleListenerRef.set(idleListener);
+ selectedMailbox.registerIdle(idleListener);
} else {
- line = "";
+ idleReadySink.tryEmitEmpty();
}
- if (sm != null) {
- sm.unregisterIdle();
+ if (!idleActive.get()) {
+ cleanupIdle(session, selectedMailbox, idleActive,
lineHandlerState, idleReadySink, idleListener);
+ 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 {
+ cleanupIdle(session, selectedMailbox, idleActive,
lineHandlerState, idleReadySink, idleListener);
+ } 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) {
+ return (session1, data) -> {
+ if (!idleActive.get()) {
+ return Mono.empty();
+ }
+ lineHandlerState.compareAndSet(LineHandlerState.INSTALLING,
LineHandlerState.INSTALLED);
Review Comment:
The Linearizer is supposed to provide sequential execution of IMAP commands
on the same session.
Given that invariant for IMAP workload, most concurrency based bugs
identified by looking at IdleProcessor alone (and reproduced through solely
unit tests) are never happening; they thus do not justify making the code more
complex.
Have the IMAP invariant of sequential command execution been taken into
account when designing this state machine? Am I missing something?
--
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]