chibenwa commented on code in PR #3198:
URL: https://github.com/apache/james-project/pull/3198#discussion_r4090532687
##########
protocols/imap/src/main/java/org/apache/james/imap/processor/IdleProcessor.java:
##########
@@ -77,56 +77,83 @@ 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))
+ SelectedMailbox sm = session.getSelected();
+ Sinks.One<Void> idleReadySink = Sinks.one();
+ AtomicBoolean idleActive = new AtomicBoolean(true);
+ AtomicBoolean lineHandlerInstalled = new AtomicBoolean(false);
+ return Mono.fromRunnable(() -> idle(request, session, responder, sm,
idleReadySink, idleActive, lineHandlerInstalled))
.then(unsolicitedResponses(session, responder, false))
.onErrorResume(e -> {
+ cleanupIdle(session, sm, idleActive, lineHandlerInstalled,
idleReadySink);
no(request, responder,
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 void cleanupIdle(ImapSession session, SelectedMailbox sm,
AtomicBoolean idleActive,
+ AtomicBoolean lineHandlerInstalled,
Sinks.One<Void> idleReadySink) {
+ if (idleActive.compareAndSet(true, false)) {
+ if (sm != null) {
+ sm.unregisterIdle();
+ }
+ if (lineHandlerInstalled.get()) {
+ session.popLineHandler();
+ }
+ idleReadySink.tryEmitEmpty();
}
+ }
- final AtomicBoolean idleActive = new AtomicBoolean(true);
+ private void idle(IdleRequest request, ImapSession session, Responder
responder, SelectedMailbox sm,
+ Sinks.One<Void> idleReadySink, AtomicBoolean idleActive,
AtomicBoolean lineHandlerInstalled) {
+ if (sm != null) {
+ sm.registerIdle(new IdleMailboxListener(session, sm, responder,
idleReadySink, idleActive, lineHandlerInstalled));
+ } else {
+ idleReadySink.tryEmitEmpty();
+ }
- session.pushLineHandler((session1, data) -> Mono.fromRunnable(() -> {
- String line;
- if (data.length > 2) {
- line = new String(data, 0, data.length - 2);
- } else {
- line = "";
- }
+ try {
+ session.pushLineHandler((session1, data) -> {
+ lineHandlerInstalled.set(true);
+ cleanupIdle(session1, sm, idleActive, lineHandlerInstalled,
idleReadySink);
+ String line = new String(data,
StandardCharsets.US_ASCII).trim();
- if (sm != null) {
- sm.unregisterIdle();
- }
- if (!DONE.equals(line.toUpperCase(Locale.US))) {
- String message = String.format("Continuation for IMAP IDLE was
not understood. Expected 'DONE', got '%s'.", line);
+ if (line.isEmpty() || !session1.isConnected()) {
+ LOGGER.debug("IDLE continuation received empty input or
disconnected session.");
+ return Mono.empty();
+ }
+ String upper = line.toUpperCase(Locale.US);
+ if (DONE.equals(upper)) {
+ okComplete(request, responder);
+ responder.flush();
+ return Mono.empty();
+ }
+ String sanitized = line.replaceAll("[\\r\\n\\x00-\\x1F]", "");
+ String displayLine = sanitized.length() > 32 ?
sanitized.substring(0, 32) + "..." : sanitized;
+ 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.info(message);
+ LOGGER.debug(message);
responder.respond(response);
responder.flush();
- } else {
- okComplete(request, responder);
- responder.flush();
- }
- session1.popLineHandler();
- idleActive.set(false);
- }));
+ return Mono.empty();
+ });
+ lineHandlerInstalled.set(true);
+ } catch (Exception e) {
+ lineHandlerInstalled.set(false);
+ throw e;
+ }
+
+ // Write the response after the listener was added (IMAP-341)
+ responder.respond(new ContinuationResponse(HumanReadableText.IDLING));
+ responder.flush();
Review Comment:
Are we sure?
By moving respond out of the line handler onto the event loop we are
answering before the line handler is set up.
Netty write ordering is a complex topic, and its behaviour depends on wether
we are running on the event loop or not.
--
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]