Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down Expand Up @@ -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()
Expand Down Expand Up @@ -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)
}
Expand Down Expand Up @@ -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 =
Comment thread
stefanmirkovic marked this conversation as resolved.
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 {
Expand Down Expand Up @@ -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
Expand All @@ -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
}
}
Expand Down
Original file line number Diff line number Diff line change
@@ -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)
}
}
Loading