KubeMQ
Client SDKsKotlinReference

Client

Client classes, DSL builders, and lifecycle -- KubeMQ Kotlin SDK reference.

The Kotlin SDK splits messaging APIs across focused clients: pub/sub (PubSubClient), queues (QueuesClient), and commands/queries (CQClient). All are created via KubeMQClient DSL builders and implement Closeable.

Package

build.gradle.kts
implementation("io.kubemq.sdk:kubemq-sdk-kotlin:1.0.1")

Requires Kotlin 2.0+ and JDK 11+.

Client classes

BuilderReturnsPatterns
KubeMQClient.pubSub { }PubSubClientEvents, Events Store
KubeMQClient.queues { }QueuesClientQueues
KubeMQClient.cq { }CQClientCommands, Queries

DSL construction

val pub = KubeMQClient.pubSub {
    address = "localhost:50000"
    clientId = "kotlin-pub"
}

val queues = KubeMQClient.queues {
    address = "localhost:50000"
    clientId = "kotlin-q"
}

val cq = KubeMQClient.cq {
    address = "localhost:50000"
    clientId = "kotlin-cq"
}

Configuration properties:

NameTypeRequiredDescription
addressStringYeshost:port for the broker
clientIdStringYesStable ID for this process
authTokenStringNoBearer token
logLevelLogLevelNoLogging level (DEBUG, INFO, WARN, ERROR)
keepAliveBooleanNoEnable gRPC keep-alive
pingIntervalSecondsIntNoKeep-alive ping interval
pingTimeoutSecondsIntNoKeep-alive ping timeout
maxReceiveSizeIntNoMax inbound message size
ensureConnectedTimeoutMsLongNoConnection timeout
unaryTimeoutMsLongNoTimeout for unary gRPC calls
tls { }DSL blockNoTLS/mTLS configuration
reconnection { }DSL blockNoReconnection backoff tuning

Reconnection DSL

reconnection {
    initialBackoffMs = 1000
    maxBackoffMs = 30_000
    multiplier = 2.0
    maxRetries = 10
}

TLS DSL

tls {
    caCertFile = "/path/to/ca.pem"
    certFile = "/path/to/client.pem"
    keyFile = "/path/to/client.key"
}

Lifecycle

Use Kotlin's use { } block (equivalent to Java try-with-resources) because each client implements Closeable:

KubeMQClient.pubSub {
    address = "localhost:50000"
    clientId = "demo"
}.use { client ->
    // use client
}

Connection state monitoring

client.connectionState.collect { state ->
    when (state) {
        is ConnectionState.Idle -> println("Idle")
        is ConnectionState.Connecting -> println("Connecting")
        is ConnectionState.Ready -> println("Ready")
        is ConnectionState.Reconnecting -> println("Reconnecting #${state.attempt}")
        is ConnectionState.Closed -> println("Closed")
    }
}

Environment configuration

val config = ClientConfig.fromEnvironment()
// Reads KUBEMQ_ADDRESS, KUBEMQ_CLIENT_ID, KUBEMQ_AUTH_TOKEN

Quick Usage

App.kt
KubeMQClient.pubSub {
    address = "localhost:50000"
    clientId = "reference-demo"
}.use { pub ->
    // Events APIs live on PubSubClient -- see Events reference page
}

See Also

Was this page helpful?

On this page