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
4 changes: 2 additions & 2 deletions app/src/main/java/com/gojek/courier/app/ui/MainActivity.kt
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ import com.gojek.courier.streamadapter.rxjava2.RxJava2StreamAdapterFactory
import com.gojek.mqtt.auth.Authenticator
import com.gojek.mqtt.client.MqttClient
import com.gojek.mqtt.client.config.ExperimentConfigs
import com.gojek.mqtt.client.config.PersistenceOptions.PahoPersistenceOptions
import com.gojek.mqtt.client.config.PersistenceOptions
import com.gojek.mqtt.client.config.v3.MqttV3Configuration
import com.gojek.mqtt.client.factory.MqttClientFactory
import com.gojek.mqtt.event.EventHandler
Expand Down Expand Up @@ -158,7 +158,7 @@ class MainActivity : AppCompatActivity() {
}
},
mqttInterceptorList = listOf(MqttChuckInterceptor(this, MqttChuckConfig(retentionPeriod = Period.ONE_HOUR))),
persistenceOptions = PahoPersistenceOptions(100, false),
persistenceOptions = PersistenceOptions(bufferCapacity = 100, isDeleteOldestMessages = false),
experimentConfigs = ExperimentConfigs(
adaptiveKeepAliveConfig = AdaptiveKeepAliveConfig(
lowerBoundMinutes = 1,
Expand Down
33 changes: 17 additions & 16 deletions mqtt-client/api/mqtt-client.api
Original file line number Diff line number Diff line change
Expand Up @@ -29,13 +29,12 @@ public abstract interface class com/gojek/mqtt/client/MqttInterceptor {

public final class com/gojek/mqtt/client/config/ExperimentConfigs {
public fun <init> ()V
public fun <init> (Lcom/gojek/mqtt/client/config/SubscriptionStore;Lcom/gojek/mqtt/model/AdaptiveKeepAliveConfig;IIIIJJZIZZZ)V
public synthetic fun <init> (Lcom/gojek/mqtt/client/config/SubscriptionStore;Lcom/gojek/mqtt/model/AdaptiveKeepAliveConfig;IIIIJJZIZZZILkotlin/jvm/internal/DefaultConstructorMarker;)V
public fun <init> (Lcom/gojek/mqtt/client/config/SubscriptionStore;Lcom/gojek/mqtt/model/AdaptiveKeepAliveConfig;IIIIJJZIZZ)V
public synthetic fun <init> (Lcom/gojek/mqtt/client/config/SubscriptionStore;Lcom/gojek/mqtt/model/AdaptiveKeepAliveConfig;IIIIJJZIZZILkotlin/jvm/internal/DefaultConstructorMarker;)V
public final fun component1 ()Lcom/gojek/mqtt/client/config/SubscriptionStore;
public final fun component10 ()I
public final fun component11 ()Z
public final fun component12 ()Z
public final fun component13 ()Z
public final fun component2 ()Lcom/gojek/mqtt/model/AdaptiveKeepAliveConfig;
public final fun component3 ()I
public final fun component4 ()I
Expand All @@ -44,8 +43,8 @@ public final class com/gojek/mqtt/client/config/ExperimentConfigs {
public final fun component7 ()J
public final fun component8 ()J
public final fun component9 ()Z
public final fun copy (Lcom/gojek/mqtt/client/config/SubscriptionStore;Lcom/gojek/mqtt/model/AdaptiveKeepAliveConfig;IIIIJJZIZZZ)Lcom/gojek/mqtt/client/config/ExperimentConfigs;
public static synthetic fun copy$default (Lcom/gojek/mqtt/client/config/ExperimentConfigs;Lcom/gojek/mqtt/client/config/SubscriptionStore;Lcom/gojek/mqtt/model/AdaptiveKeepAliveConfig;IIIIJJZIZZZILjava/lang/Object;)Lcom/gojek/mqtt/client/config/ExperimentConfigs;
public final fun copy (Lcom/gojek/mqtt/client/config/SubscriptionStore;Lcom/gojek/mqtt/model/AdaptiveKeepAliveConfig;IIIIJJZIZZ)Lcom/gojek/mqtt/client/config/ExperimentConfigs;
public static synthetic fun copy$default (Lcom/gojek/mqtt/client/config/ExperimentConfigs;Lcom/gojek/mqtt/client/config/SubscriptionStore;Lcom/gojek/mqtt/model/AdaptiveKeepAliveConfig;IIIIJJZIZZILjava/lang/Object;)Lcom/gojek/mqtt/client/config/ExperimentConfigs;
public fun equals (Ljava/lang/Object;)Z
public final fun getActivityCheckIntervalSeconds ()I
public final fun getAdaptiveKeepAliveConfig ()Lcom/gojek/mqtt/model/AdaptiveKeepAliveConfig;
Expand All @@ -56,7 +55,6 @@ public final class com/gojek/mqtt/client/config/ExperimentConfigs {
public final fun getMaxInflightMessagesLimit ()I
public final fun getPolicyResetTimeSeconds ()I
public final fun getShouldSendMessageViaHandler ()Z
public final fun getShouldUseMemoryPersistence ()Z
public final fun getShouldUseNewSSLFlow ()Z
public final fun getStopMqttThreadOnDestroy ()Z
public final fun getSubscriptionStore ()Lcom/gojek/mqtt/client/config/SubscriptionStore;
Expand All @@ -80,21 +78,24 @@ public abstract class com/gojek/mqtt/client/config/MqttConfiguration {
public fun getWakeLockTimeout ()I
}

public abstract class com/gojek/mqtt/client/config/PersistenceOptions {
}

public final class com/gojek/mqtt/client/config/PersistenceOptions$PahoPersistenceOptions : com/gojek/mqtt/client/config/PersistenceOptions {
public final class com/gojek/mqtt/client/config/PersistenceOptions {
public fun <init> ()V
public fun <init> (IZ)V
public synthetic fun <init> (IZILkotlin/jvm/internal/DefaultConstructorMarker;)V
public final fun component1 ()I
public final fun component2 ()Z
public final fun copy (IZ)Lcom/gojek/mqtt/client/config/PersistenceOptions$PahoPersistenceOptions;
public static synthetic fun copy$default (Lcom/gojek/mqtt/client/config/PersistenceOptions$PahoPersistenceOptions;IZILjava/lang/Object;)Lcom/gojek/mqtt/client/config/PersistenceOptions$PahoPersistenceOptions;
public fun <init> (ZIIZZ)V
public synthetic fun <init> (ZIIZZILkotlin/jvm/internal/DefaultConstructorMarker;)V
public final fun component1 ()Z
public final fun component2 ()I
public final fun component3 ()I
public final fun component4 ()Z
public final fun component5 ()Z
public final fun copy (ZIIZZ)Lcom/gojek/mqtt/client/config/PersistenceOptions;
public static synthetic fun copy$default (Lcom/gojek/mqtt/client/config/PersistenceOptions;ZIIZZILjava/lang/Object;)Lcom/gojek/mqtt/client/config/PersistenceOptions;
public fun equals (Ljava/lang/Object;)Z
public final fun getBufferCapacity ()I
public final fun getMemoryPersistenceCapacity ()I
public final fun getShouldUseMemoryPersistence ()Z
public fun hashCode ()I
public final fun isDeleteOldestMessages ()Z
public final fun isPersistBuffer ()Z
public fun toString ()Ljava/lang/String;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,6 @@ data class ExperimentConfigs(
val shouldUseNewSSLFlow: Boolean = false,
val maxInflightMessagesLimit: Int = MAX_INFLIGHT_MESSAGES_ALLOWED,
val stopMqttThreadOnDestroy: Boolean = false,
val shouldUseMemoryPersistence: Boolean = false,
val shouldSendMessageViaHandler: Boolean = true
)

Expand Down
Original file line number Diff line number Diff line change
@@ -1,11 +1,13 @@
package com.gojek.mqtt.client.config

sealed class PersistenceOptions {
data class PahoPersistenceOptions(
val bufferCapacity: Int = OFFLINE_BUFFER_CAPACITY_DEFAULT,
val isDeleteOldestMessages: Boolean = DELETE_OLDEST_MESSAGES_DEFAULT
) : PersistenceOptions()
}
data class PersistenceOptions(
val shouldUseMemoryPersistence: Boolean = false,
val memoryPersistenceCapacity: Int = 100,
val bufferCapacity: Int = OFFLINE_BUFFER_CAPACITY_DEFAULT,
val isPersistBuffer: Boolean = PERSIST_BUFFER_DEFAULT,
val isDeleteOldestMessages: Boolean = DELETE_OLDEST_MESSAGES_DEFAULT
)

private const val OFFLINE_BUFFER_CAPACITY_DEFAULT = 50000
private const val PERSIST_BUFFER_DEFAULT = true
private const val DELETE_OLDEST_MESSAGES_DEFAULT = false
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,6 @@ import com.gojek.mqtt.client.MqttInterceptor
import com.gojek.mqtt.client.config.ExperimentConfigs
import com.gojek.mqtt.client.config.MqttConfiguration
import com.gojek.mqtt.client.config.PersistenceOptions
import com.gojek.mqtt.client.config.PersistenceOptions.PahoPersistenceOptions
import com.gojek.mqtt.constants.DEFAULT_WAKELOCK_TIMEOUT
import com.gojek.mqtt.exception.handler.v3.AuthFailureHandler
import com.gojek.mqtt.pingsender.MqttPingSender
Expand Down Expand Up @@ -36,7 +35,7 @@ data class MqttV3Configuration(
override val authFailureHandler: AuthFailureHandler? = null,
override val pingSender: MqttPingSender,
override val mqttInterceptorList: List<MqttInterceptor> = emptyList(),
override val persistenceOptions: PersistenceOptions = PahoPersistenceOptions(),
override val persistenceOptions: PersistenceOptions = PersistenceOptions(),
override val experimentConfigs: ExperimentConfigs = ExperimentConfigs()
) : MqttConfiguration(
connectRetryTimePolicy = connectRetryTimePolicy,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -162,8 +162,7 @@ internal class AndroidMqttClient(
inactivityTimeoutSeconds = experimentConfigs.inactivityTimeoutSeconds,
policyResetTimeSeconds = experimentConfigs.policyResetTimeSeconds,
shouldUseNewSSLFlow = experimentConfigs.shouldUseNewSSLFlow,
connectPacketTimeoutSeconds = experimentConfigs.connectPacketTimeoutSeconds,
shouldUseMemoryPersistence = experimentConfigs.shouldUseMemoryPersistence
connectPacketTimeoutSeconds = experimentConfigs.connectPacketTimeoutSeconds
)

mqttConnection = MqttConnection(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,6 @@ import com.gojek.courier.logging.ILogger
import com.gojek.courier.utils.Clock
import com.gojek.keepalive.KeepAliveFailureHandler
import com.gojek.mqtt.client.IMessageReceiveListener
import com.gojek.mqtt.client.config.PersistenceOptions.PahoPersistenceOptions
import com.gojek.mqtt.client.model.MqttSendPacket
import com.gojek.mqtt.connection.config.v3.ConnectionConfig
import com.gojek.mqtt.event.PahoEventHandler
Expand Down Expand Up @@ -49,6 +48,7 @@ import org.eclipse.paho.client.mqttv3.MqttSecurityException
import org.eclipse.paho.client.mqttv3.internal.wire.MqttSuback
import org.eclipse.paho.client.mqttv3.internal.wire.SubscribeFlags
import org.eclipse.paho.client.mqttv3.internal.wire.UserProperty
import org.eclipse.paho.client.mqttv3.persist.BoundedMemoryPersistence
import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence

internal class MqttConnection(
Expand Down Expand Up @@ -402,8 +402,16 @@ internal class MqttConnection(
}

private fun getMqttAsyncClient(clientId: String, serverUri: String): MqttAsyncClient {
val persistence = if (connectionConfig.shouldUseMemoryPersistence) {
MemoryPersistence()
val persistenceOptions = connectionConfig.persistenceOptions
val persistence = if (persistenceOptions.shouldUseMemoryPersistence) {
if (persistenceOptions.memoryPersistenceCapacity > 0) {
BoundedMemoryPersistence(
persistenceOptions.memoryPersistenceCapacity,
persistenceOptions.isDeleteOldestMessages
)
} else {
MemoryPersistence()
}
} else {
pahoPersistence
}
Expand All @@ -420,9 +428,9 @@ internal class MqttConnection(
connectionConfig.mqttInterceptorList
)
val bufferOptions = DisconnectedBufferOptions()
with(connectionConfig.persistenceOptions as PahoPersistenceOptions) {
with(connectionConfig.persistenceOptions) {
bufferOptions.isBufferEnabled = true
bufferOptions.isPersistBuffer = true
bufferOptions.isPersistBuffer = isPersistBuffer
bufferOptions.bufferSize = bufferCapacity
bufferOptions.isDeleteOldestMessages = isDeleteOldestMessages
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,5 @@ internal data class ConnectionConfig(
val inactivityTimeoutSeconds: Int,
val policyResetTimeSeconds: Int,
val shouldUseNewSSLFlow: Boolean,
val shouldUseMemoryPersistence: Boolean,
val connectPacketTimeoutSeconds: Int
)
Original file line number Diff line number Diff line change
@@ -0,0 +1,114 @@
/*******************************************************************************
* Copyright (c) 2009, 2014 IBM Corp.
*
* All rights reserved. This program and the accompanying materials
* are made available under the terms of the Eclipse Public License v1.0
* and Eclipse Distribution License v1.0 which accompany this distribution.
*
* The Eclipse Public License is available at
* http://www.eclipse.org/legal/epl-v10.html
* and the Eclipse Distribution License is available at
* http://www.eclipse.org/org/documents/edl-v10.php.
*
* Contributors:
* Dave Locke - initial API and implementation and/or initial documentation
*/
package org.eclipse.paho.client.mqttv3.persist;

import org.eclipse.paho.client.mqttv3.MqttClientPersistence;
import org.eclipse.paho.client.mqttv3.MqttPersistable;
import org.eclipse.paho.client.mqttv3.MqttPersistenceException;

import java.util.Collections;
import java.util.Enumeration;
import java.util.LinkedHashMap;
import java.util.Map;

/**
* Persistence that uses memory
*
* In cases where reliability is not required across client or device
* restarts memory this memory peristence can be used. In cases where
* reliability is required like when clean session is set to false
* then a non-volatile form of persistence should be used.
*
*/
public class BoundedMemoryPersistence implements MqttClientPersistence {

private Map<String, MqttPersistable> data;
private final int maxCapacity;
private final boolean isDeleteOldestMessages;

public BoundedMemoryPersistence() {
this(100, false);
}

public BoundedMemoryPersistence(int maxCapacity, boolean isDeleteOldestMessages) {
this.maxCapacity = maxCapacity;
this.isDeleteOldestMessages = isDeleteOldestMessages;
}

/* (non-Javadoc)
* @see org.eclipse.paho.client.mqttv3.MqttClientPersistence#close()
*/
public void close() throws MqttPersistenceException {
data.clear();
}

/* (non-Javadoc)
* @see org.eclipse.paho.client.mqttv3.MqttClientPersistence#keys()
*/
public Enumeration keys() throws MqttPersistenceException {
return Collections.enumeration(data.keySet());
}

/* (non-Javadoc)
* @see org.eclipse.paho.client.mqttv3.MqttClientPersistence#get(java.lang.String)
*/
public MqttPersistable get(String key) throws MqttPersistenceException {
return (MqttPersistable)data.get(key);
}

/* (non-Javadoc)
* @see org.eclipse.paho.client.mqttv3.MqttClientPersistence#open(java.lang.String, java.lang.String)
*/
public void open(String clientId, String serverURI) throws MqttPersistenceException {
this.data = Collections.synchronizedMap(new LinkedHashMap<String, MqttPersistable>());
}

/* (non-Javadoc)
* @see org.eclipse.paho.client.mqttv3.MqttClientPersistence#put(java.lang.String, org.eclipse.paho.client.mqttv3.MqttPersistable)
*/
public void put(String key, MqttPersistable persistable) throws MqttPersistenceException {
synchronized (data) {
if (data.size() < maxCapacity || data.containsKey(key)) {
data.put(key, persistable);
} else if (isDeleteOldestMessages) {
String oldestKey = data.keySet().iterator().next();
data.remove(oldestKey);
data.put(key, persistable);
}
}
}

/* (non-Javadoc)
* @see org.eclipse.paho.client.mqttv3.MqttClientPersistence#remove(java.lang.String)
*/
public void remove(String key) throws MqttPersistenceException {
data.remove(key);
}

/* (non-Javadoc)
* @see org.eclipse.paho.client.mqttv3.MqttClientPersistence#clear()
*/
public void clear() throws MqttPersistenceException {
data.clear();
}

/* (non-Javadoc)
* @see org.eclipse.paho.client.mqttv3.MqttClientPersistence#containsKey(java.lang.String)
*/
public boolean containsKey(String key) throws MqttPersistenceException {
return data.containsKey(key);
}
}
Loading