diff --git a/framework-client/src/main/java/io/axoniq/platform/framework/client/AxoniqConsoleRSocketClient.kt b/framework-client/src/main/java/io/axoniq/platform/framework/client/AxoniqConsoleRSocketClient.kt index 3082f0b..9ae4c32 100644 --- a/framework-client/src/main/java/io/axoniq/platform/framework/client/AxoniqConsoleRSocketClient.kt +++ b/framework-client/src/main/java/io/axoniq/platform/framework/client/AxoniqConsoleRSocketClient.kt @@ -94,6 +94,9 @@ class AxoniqConsoleRSocketClient( private var status: ClientStatus = ClientStatus.PENDING private var suppressConnectMessage = false + /** Whether a connection has ever succeeded, which decides how loudly a refusal is reported. */ + private var hasEverConnected = false + init { platformClientConnectionService.subscribeToSettings(heartbeatOrchestrator) @@ -196,13 +199,18 @@ class AxoniqConsoleRSocketClient( logger.info("Connection to Axoniq Platform set up successfully! This instance's name: $instanceName, settings: $settings") suppressConnectMessage = true } + hasEverConnected = true connectionRetryCount = 0 socket } } .doOnError { e -> disposeCurrentConnection() - platformClientConnectionService.notifyUnreachable(classifyConnectionError(e)) + val reason = classifyConnectionError(e) + val refused = reason == PlatformClientConnectionObserver.UnreachableReason.INVALID_AUTHENTICATION + if (!refused || shouldReport(true, hasEverConnected, connectionRetryCount)) { + platformClientConnectionService.notifyUnreachable(reason) + } } .doFinally { synchronized(connectionLock) { pendingConnection = null } } .cache() @@ -244,15 +252,16 @@ class AxoniqConsoleRSocketClient( getOrConnectRSocket().subscribe( { /* success — logged inside buildConnectionMono */ }, { e -> - if (retryCount == 4) { - if (suppressConnectMessage) { - logger.warn("Lost connection to Axoniq Platform. Will keep trying to reconnect...") + val refusedCredentials = isAuthenticationFailure(e) + if (shouldReport(refusedCredentials, hasEverConnected, retryCount)) { + if (refusedCredentials) { + logger.info("Axoniq Platform refused this application's credentials. Check the access token and the environment it belongs to. Will keep trying to connect.") + } else if (suppressConnectMessage) { + logger.info("Lost connection to Axoniq Platform. Will keep trying to reconnect...") } else { - logger.warn("Unable to connect to Axoniq Platform. Will keep trying to reconnect...") + logger.info("Unable to connect to Axoniq Platform. Will keep trying to reconnect...") } suppressConnectMessage = false - } else if (retryCount > 4 && retryCount % 10 == 0) { - logger.error("Still unable to reconnect to Axoniq Platform after $retryCount attempts. Reason: ${e.message}") } logger.debug("Failed to connect to Axoniq Platform", e) } @@ -337,6 +346,41 @@ class AxoniqConsoleRSocketClient( companion object { private const val BACKOFF_FACTOR = 2.0 + + /** + * How often a persistent connection problem is repeated. With [BACKOFF_FACTOR] backing off to a + * minute, the first report lands a few minutes in and roughly every ten minutes after that: long + * enough for a platform-side blip to come and go unnoticed, often enough that a lasting problem + * cannot be missed. + */ + private const val RETRIES_BETWEEN_REPORTS = 10 + + private const val MAX_CAUSE_DEPTH = 5 + private val AUTH_FAILURE_MARKERS = listOf("Access Denied", "invalid authentication") + + /** + * Whether the platform refused our credentials, as opposed to being unreachable. RSocket wraps the + * server's error, so the whole cause chain is considered. + */ + internal fun isAuthenticationFailure(error: Throwable): Boolean = + generateSequence(error as Throwable?) { it.cause } + .take(MAX_CAUSE_DEPTH) + .mapNotNull { it.message } + .any { message -> AUTH_FAILURE_MARKERS.any { message.contains(it, ignoreCase = true) } } + + /** + * Whether this failure is worth a log line. + * + * A refusal on the very first attempt an application ever makes is said immediately. Nothing has + * ever worked, so a misconfigured token is far likelier than the platform having a moment, and + * whoever is starting the application is usually watching. Once a connection has been established + * the same refusal is most likely transient, and saying so at once would be crying wolf, so it + * waits for the periodic report. + */ + internal fun shouldReport(refusedCredentials: Boolean, hasEverConnected: Boolean, retryCount: Int): Boolean { + val firstEverAttempt = refusedCredentials && !hasEverConnected && retryCount == 1 + return firstEverAttempt || (retryCount > 0 && retryCount % RETRIES_BETWEEN_REPORTS == 0) + } } private inner class HeartbeatOrchestrator : PlatformClientConnectionObserver { @@ -397,9 +441,7 @@ class AxoniqConsoleRSocketClient( private fun classifyConnectionError(e: Throwable): PlatformClientConnectionObserver.UnreachableReason { return when { - e.message?.contains("invalid authentication", ignoreCase = true) == true -> - PlatformClientConnectionObserver.UnreachableReason.INVALID_AUTHENTICATION - e.message?.contains("Access Denied", ignoreCase = true) == true -> + isAuthenticationFailure(e) -> PlatformClientConnectionObserver.UnreachableReason.INVALID_AUTHENTICATION e is java.net.ConnectException || e.cause is java.net.ConnectException -> PlatformClientConnectionObserver.UnreachableReason.NO_CONNECTION @@ -414,12 +456,13 @@ class AxoniqConsoleRSocketClient( .map { encodingStrategy.decode(it, ClientSettingsV2::class.java) } + // Mapped, not logged: connectSafely decides whether a refusal has persisted long enough + // to be worth reporting. The platform can refuse a perfectly good token for a moment while + // it is itself reconnecting, and the retry loop recovers from that unaided. .onErrorMap { if (it.message?.contains("Access Denied") == true) { - logger.error("Was unable to connect to Axoniq Platform due to invalid authentication! Make sure the access token is correct.") IllegalStateException("Was unable to connect to Axoniq Platform due to invalid authentication! Make sure the access token is correct.") } else { - logger.error("Could not connect to Axoniq Platform due to connection error: ${it.message}", it) it } } diff --git a/framework-client/src/test/kotlin/io/axoniq/platform/framework/client/AxoniqConsoleRSocketClientLoggingTest.kt b/framework-client/src/test/kotlin/io/axoniq/platform/framework/client/AxoniqConsoleRSocketClientLoggingTest.kt new file mode 100644 index 0000000..2f4776f --- /dev/null +++ b/framework-client/src/test/kotlin/io/axoniq/platform/framework/client/AxoniqConsoleRSocketClientLoggingTest.kt @@ -0,0 +1,130 @@ +/* + * Copyright (c) 2026. AxonIQ B.V. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package io.axoniq.platform.framework.client + +import org.junit.jupiter.api.Nested +import org.junit.jupiter.api.Test +import kotlin.test.assertFalse +import kotlin.test.assertTrue + +class AxoniqConsoleRSocketClientLoggingTest { + + @Nested + inner class RecognisingARefusal { + + @Test + fun `recognises what the platform says when it rejects a token`() { + assertTrue(isAuthFailure(RuntimeException("Access Denied"))) + assertTrue(isAuthFailure(RuntimeException("access denied"))) + assertTrue(isAuthFailure(RuntimeException("invalid authentication"))) + } + + @Test + fun `looks through the wrapper RSocket puts around the server's error`() { + val wrapped = IllegalStateException( + "Could not receive the settings from Axoniq Platform!", + RuntimeException("Access Denied") + ) + + assertTrue(isAuthFailure(wrapped)) + } + + @Test + fun `does not mistake an unreachable platform for a rejected token`() { + assertFalse(isAuthFailure(java.net.ConnectException("Connection refused"))) + assertFalse(isAuthFailure(RuntimeException("Connection reset by peer"))) + assertFalse(isAuthFailure(RuntimeException(null as String?))) + } + + @Test + fun `does not mistake a network failure that merely mentions authentication`() { + // Telling an operator to check their access token because a proxy or TLS handshake said + // "authentication" is the misattribution this all exists to remove. + assertFalse(isAuthFailure(RuntimeException("Proxy Authentication Required"))) + assertFalse(isAuthFailure(RuntimeException("SSL handshake failed: client authentication"))) + assertFalse(isAuthFailure(RuntimeException("Unauthorized"))) + } + + @Test + fun `terminates on a cause chain that refers back to itself`() { + val looping = object : RuntimeException("Connection reset") { + override val cause: Throwable get() = this + } + + assertFalse(isAuthFailure(looping)) + } + + private fun isAuthFailure(error: Throwable) = + AxoniqConsoleRSocketClient.isAuthenticationFailure(error) + } + + @Nested + inner class DecidingWhenToSpeak { + + @Test + fun `says so at once when the very first connection an application makes is refused`() { + // Nothing has ever worked, so this is far likelier to be a misconfigured token than a blip, + // and someone is usually watching the application start. + assertTrue(shouldReport(refusedCredentials = true, hasEverConnected = false, retryCount = 1)) + } + + @Test + fun `stays quiet when a refusal interrupts a connection that had been working`() { + // The case the customer hit: the platform refused a perfectly good token for two seconds while + // its own gateway was restarting. + assertFalse(shouldReport(refusedCredentials = true, hasEverConnected = true, retryCount = 1)) + (2..9).forEach { + assertFalse(shouldReport(refusedCredentials = true, hasEverConnected = true, retryCount = it)) + } + } + + @Test + fun `repeats itself while the problem lasts`() { + listOf(10, 20, 30).forEach { + assertTrue(shouldReport(refusedCredentials = true, hasEverConnected = true, retryCount = it)) + assertTrue(shouldReport(refusedCredentials = false, hasEverConnected = true, retryCount = it)) + } + listOf(11, 19, 21).forEach { + assertFalse(shouldReport(refusedCredentials = true, hasEverConnected = true, retryCount = it)) + } + } + + @Test + fun `keeps reporting a token that never worked, after the immediate first word`() { + assertTrue(shouldReport(refusedCredentials = true, hasEverConnected = false, retryCount = 1)) + (2..9).forEach { + assertFalse(shouldReport(refusedCredentials = true, hasEverConnected = false, retryCount = it)) + } + assertTrue(shouldReport(refusedCredentials = true, hasEverConnected = false, retryCount = 10)) + } + + @Test + fun `does not treat an unreachable platform at startup as something to shout about`() { + // A platform that cannot be reached on the first attempt usually can be on the second. + assertFalse(shouldReport(refusedCredentials = false, hasEverConnected = false, retryCount = 1)) + } + + @Test + fun `never reports before an attempt has been made`() { + assertFalse(shouldReport(refusedCredentials = true, hasEverConnected = false, retryCount = 0)) + assertFalse(shouldReport(refusedCredentials = false, hasEverConnected = true, retryCount = 0)) + } + + private fun shouldReport(refusedCredentials: Boolean, hasEverConnected: Boolean, retryCount: Int) = + AxoniqConsoleRSocketClient.shouldReport(refusedCredentials, hasEverConnected, retryCount) + } +}