From 23cf38445f4a75b7cff598dc3f39fd1b34c04562 Mon Sep 17 00:00:00 2001 From: Pritimay Sarkar Date: Tue, 13 Feb 2024 08:59:03 +0530 Subject: [PATCH] handle method and dynamic config for NATS --- .../hpostesting/data/dao/MyDataBase.kt | 2 +- .../data/model/patient/DeviceData.kt | 8 +- .../hpostesting/presentation/NatsManager.kt | 194 ++++++++++-------- .../assurance/AssuranceControlsActivity.kt | 23 +-- .../dashboard/DashboardActivity.kt | 32 ++- 5 files changed, 139 insertions(+), 120 deletions(-) diff --git a/app/src/main/java/com/example/hpostesting/data/dao/MyDataBase.kt b/app/src/main/java/com/example/hpostesting/data/dao/MyDataBase.kt index 95a8940..fcfc0cc 100644 --- a/app/src/main/java/com/example/hpostesting/data/dao/MyDataBase.kt +++ b/app/src/main/java/com/example/hpostesting/data/dao/MyDataBase.kt @@ -8,7 +8,7 @@ import com.example.hpostesting.data.model.patient.DeviceData import com.example.hpostesting.data.model.patient.HemoCubeTestData import com.example.hpostesting.data.model.patient.UserData -@Database(entities = [UserData::class, HemoCubeTestData::class, DeviceData::class, BufferCheckData::class], version = 25, exportSchema = false) +@Database(entities = [UserData::class, HemoCubeTestData::class, DeviceData::class, BufferCheckData::class], version = 26, exportSchema = false) @TypeConverters(Converters::class) abstract class MyDatabase : RoomDatabase() { abstract fun userDao(): UserDao diff --git a/app/src/main/java/com/example/hpostesting/data/model/patient/DeviceData.kt b/app/src/main/java/com/example/hpostesting/data/model/patient/DeviceData.kt index 5ccb755..79b3834 100644 --- a/app/src/main/java/com/example/hpostesting/data/model/patient/DeviceData.kt +++ b/app/src/main/java/com/example/hpostesting/data/model/patient/DeviceData.kt @@ -2,7 +2,6 @@ package com.example.hpostesting.data.model.patient import androidx.room.Entity import androidx.room.PrimaryKey -import com.example.hpostesting.data.model.deviceprovision.DeviceProvisionResponse import com.google.firebase.firestore.PropertyName @Entity(tableName = "device_table") @@ -22,6 +21,9 @@ data class DeviceData( @get:PropertyName("username") @set:PropertyName("username") var username: String = "", @get:PropertyName("password") @set:PropertyName("password") - var password: String = "" - + var password: String = "", + @get:PropertyName("natsToken") @set:PropertyName("natsToken") + var natsToken: String = "", + @get:PropertyName("natsTokenExpiry") @set:PropertyName("natsTokenExpiry") + var natsTokenExpiry: String = "" ) \ No newline at end of file diff --git a/app/src/main/java/com/example/hpostesting/presentation/NatsManager.kt b/app/src/main/java/com/example/hpostesting/presentation/NatsManager.kt index ccc8b51..36cd967 100644 --- a/app/src/main/java/com/example/hpostesting/presentation/NatsManager.kt +++ b/app/src/main/java/com/example/hpostesting/presentation/NatsManager.kt @@ -42,7 +42,8 @@ class NatsManager(datacollector: DashboardActivity) { FileInputStream(clientCertPath).use { keyStoreInputStream -> keyStore.load(keyStoreInputStream, keyStorePassword) } - val caCertPath = "/storage/sdcard0/Android/data/in.sminnovations.hpostesting.quality/files/NATS/clientCertificate/client-cert.pem" + val caCertPath = + "/storage/sdcard0/Android/data/in.sminnovations.hpostesting.quality/files/NATS/clientCertificate/client-cert.pem" val caCert = FileInputStream(caCertPath).use { inputStream -> val certificateFactory = CertificateFactory.getInstance("X.509") certificateFactory.generateCertificate(inputStream) @@ -73,47 +74,50 @@ class NatsManager(datacollector: DashboardActivity) { Log.d(TAG, "TRY TO CONNECT") Thread { - val seedString = sharedPreferences.getString(Constants.NATS_TOKEN, "") - Log.e("seedString", seedString.toString()) - val seedBytes = seedString?.toCharArray() - - val theNKey = NKey.fromSeed(seedBytes) // really should load from somewhere - - val options = Options.Builder() - .server("nats://nanodgx.in:4222") - .sslContext(SSLUtils.createOpenTLSContext()) - .authHandler(object : AuthHandler { - override fun getID(): CharArray? { - return try { - theNKey?.publicKey - } catch (ex: GeneralSecurityException) { - null - } catch (ex: IOException) { - null - } catch (ex: NullPointerException) { - null - } - } - - override fun sign(nonce: ByteArray): ByteArray? { - return try { - theNKey?.sign(nonce) - } catch (ex: GeneralSecurityException) { - null - } catch (ex: IOException) { - null - } catch (ex: NullPointerException) { - null - } - } - - override fun getJWT(): CharArray? { - return null - } - }) - .build() - try { + + val seedString = sharedPreferences.getString(Constants.NATS_TOKEN, "") + val deviceId = sharedPreferences.getString(Constants.DEVICE_ID, "") + Log.e("seedString", seedString.toString()) + Log.d("nats deviceId", deviceId.toString()) + val seedBytes = seedString?.toCharArray() + + val theNKey = NKey.fromSeed(seedBytes) // really should load from somewhere + + val options = Options.Builder() + .server("nats://nanodgx.in:4222") + .sslContext(SSLUtils.createOpenTLSContext()) + .authHandler(object : AuthHandler { + override fun getID(): CharArray? { + return try { + theNKey?.publicKey + } catch (ex: GeneralSecurityException) { + null + } catch (ex: IOException) { + null + } catch (ex: NullPointerException) { + null + } + } + + override fun sign(nonce: ByteArray): ByteArray? { + return try { + theNKey?.sign(nonce) + } catch (ex: GeneralSecurityException) { + null + } catch (ex: IOException) { + null + } catch (ex: NullPointerException) { + null + } + } + + override fun getJWT(): CharArray? { + return null + } + }) + .build() + nc = Nats.connect(options) Log.d(TAG, "Connected to Nats server ${options.servers.first()}") connect = true @@ -121,56 +125,62 @@ class NatsManager(datacollector: DashboardActivity) { if (nc?.status == Connection.Status.CONNECTED) { Log.d("NATSCONNECTION", "NATS is successfully connected.") + + val d = nc?.createDispatcher { msg: Message? -> + println("Nats dispatcher $msg") + } + + nc?.subscribe("device.hpos.${deviceId}.ping") + + nc?.publish( + "server.hpos.${deviceId}.ping", + "ALIVE".toByteArray(StandardCharsets.UTF_8) + ) + nc?.publish( + "server.hpos.${deviceId}.health", + "ALIVE".toByteArray(StandardCharsets.UTF_8) + ) + + d?.subscribe("device.hpos.${deviceId}.ping") { msg -> + val response = String(msg.data, StandardCharsets.UTF_8) + datacollector.setResponse(response) + println("Message received (up to 100 times): $response") + Log.d(TAG, "subscribed msg ${msg} on topic ping") + } + + d?.subscribe("device.hpos.${deviceId}.update") { msg -> + val response = String(msg.data, StandardCharsets.UTF_8) + datacollector.setResponse(response) + println("Message received (up to 100 times): $response") + } + + d?.subscribe("device.hpos.${deviceId}.uploadlogs") { msg -> + val response = String(msg.data, StandardCharsets.UTF_8) + datacollector.setResponse(response + "UPLOAD") + println("Message received (up to 100 times): $response") + } + + d?.subscribe("device.hpos.${deviceId}.disable") { msg -> + val response = String(msg.data, StandardCharsets.UTF_8) + datacollector.setResponse(response) + println("Message received (up to 100 times): $response") + } + + d?.subscribe("device.hpos.${deviceId}.updatecustomer") { msg -> + val response = String(msg.data, StandardCharsets.UTF_8) + datacollector.setResponse(response) + println("Message received (up to 100 times): $response") + } + + d?.subscribe("device.hpos.${deviceId}.checkupdate") { msg -> + val response = String(msg.data, StandardCharsets.UTF_8) + datacollector.setResponse(response) + println("Message received (up to 100 times): $response") + } } else { Log.d("NATSCONNECTION", "NATS is not connected. Current status: ${nc?.status}") } - val deviceId = sharedPreferences.getString(Constants.DEVICE_ID, "") - - nc?.publish( - "server.hpos.${deviceId}.ping", - "ALIVE".toByteArray(StandardCharsets.UTF_8) - ) - nc?.publish( - "server.hpos.${deviceId}.health", - "ALIVE".toByteArray(StandardCharsets.UTF_8) - ) - - val d = nc?.createDispatcher { msg: Message? -> - println("Nats dispatcher $msg") - } - - - d?.subscribe("device.hpos.${deviceId}.update") { msg -> - val response = String(msg.data, StandardCharsets.UTF_8) - datacollector.setResponse(response) - println("Message received (up to 100 times): $response") - } - - d?.subscribe("device.hpos.${deviceId}.uploadlogs") { msg -> - val response = String(msg.data, StandardCharsets.UTF_8) - datacollector.setResponse(response + "UPLOAD") - println("Message received (up to 100 times): $response") - } - - d?.subscribe("device.hpos.${deviceId}.disable") { msg -> - val response = String(msg.data, StandardCharsets.UTF_8) - datacollector.setResponse(response) - println("Message received (up to 100 times): $response") - } - - d?.subscribe("device.hpos.${deviceId}.updatecustomer") { msg -> - val response = String(msg.data, StandardCharsets.UTF_8) - datacollector.setResponse(response) - println("Message received (up to 100 times): $response") - } - - d?.subscribe("device.hpos.${deviceId}.checkupdate") { msg -> - val response = String(msg.data, StandardCharsets.UTF_8) - datacollector.setResponse(response) - println("Message received (up to 100 times): $response") - } - } catch (exp: Exception) { println(exp.printStackTrace()) connect = false @@ -185,6 +195,16 @@ class NatsManager(datacollector: DashboardActivity) { Log.d(TAG, "Published msg ${msg} on topic ${topic}") } + fun sub(topic: String) { + val d = nc?.createDispatcher { msg: Message? -> + val response = String(msg?.data ?: ByteArray(0), StandardCharsets.UTF_8) + datacollector.onMessageReceived(topic, response) + Log.d(TAG, "Subscribed msg $msg on topic $topic") + } + + d?.subscribe(topic) + } + fun close() { nc?.close() Log.d(TAG, "Nats connection close") diff --git a/app/src/main/java/com/example/hpostesting/presentation/assurance/AssuranceControlsActivity.kt b/app/src/main/java/com/example/hpostesting/presentation/assurance/AssuranceControlsActivity.kt index 027fe09..307ac08 100644 --- a/app/src/main/java/com/example/hpostesting/presentation/assurance/AssuranceControlsActivity.kt +++ b/app/src/main/java/com/example/hpostesting/presentation/assurance/AssuranceControlsActivity.kt @@ -3,19 +3,14 @@ package com.example.hpostesting.presentation.assurance import android.content.Context import android.content.SharedPreferences import android.os.Bundle -import android.util.Log import androidx.appcompat.app.AppCompatActivity -import com.example.hpostesting.presentation.NatsManager -import com.example.hpostesting.presentation.dashboard.IDataCollector import dagger.hilt.android.AndroidEntryPoint import `in`.sminnovations.hpostesting.databinding.ActivityAssuranceControlsBinding @AndroidEntryPoint -class AssuranceControlsActivity: AppCompatActivity(), IDataCollector { +class AssuranceControlsActivity: AppCompatActivity() { lateinit var binding: ActivityAssuranceControlsBinding lateinit var sharedPreference: SharedPreferences - lateinit var nats: NatsManager - var responses: String = "" override fun onCreate(savedInstanceState: Bundle?) { super.onCreate(savedInstanceState) @@ -29,21 +24,5 @@ class AssuranceControlsActivity: AppCompatActivity(), IDataCollector { .replace(binding.fgAssuranceControls.id, AssuranceControlsFragment()) .commit() } - -// nats = NatsManager(this) -// nats.connect() -// nats.pub("server.hpos.HCV-000-3001.ping", "THIS IS A TEST MSG") - } - - override fun setConnect(connect: Boolean) { - if(connect){ - Log.i("NATS Connection", connect.toString()) - } - } - - override fun setResponse(response: String) { - - responses = responses+response+"\n" - println(responses) } } \ No newline at end of file diff --git a/app/src/main/java/com/example/hpostesting/presentation/dashboard/DashboardActivity.kt b/app/src/main/java/com/example/hpostesting/presentation/dashboard/DashboardActivity.kt index c4dbea1..62f1f86 100644 --- a/app/src/main/java/com/example/hpostesting/presentation/dashboard/DashboardActivity.kt +++ b/app/src/main/java/com/example/hpostesting/presentation/dashboard/DashboardActivity.kt @@ -4,7 +4,9 @@ import android.app.DownloadManager import android.content.BroadcastReceiver import android.content.Context import android.content.Intent +import android.content.SharedPreferences import android.net.Uri +import android.os.Build import android.os.Bundle import android.util.Log import android.view.Menu @@ -19,13 +21,10 @@ import androidx.navigation.ui.navigateUp import androidx.navigation.ui.setupActionBarWithNavController import androidx.navigation.ui.setupWithNavController import com.example.hpostesting.data.Result -import com.example.hpostesting.data.api.DeviceCommunicationHandler -import com.example.hpostesting.data.constant.HemoCubeCommands +import com.example.hpostesting.data.constant.Constants import com.example.hpostesting.data.constant.LanguageManager import com.example.hpostesting.presentation.NatsManager -import com.example.hpostesting.presentation.UsbServiceListener import com.example.hpostesting.presentation.hemocube.HemoCubeViewModel -import com.example.hpostesting.presentation.testRight.UsbService import com.google.android.material.navigation.NavigationView import com.google.firebase.appdistribution.FirebaseAppDistribution import com.google.firebase.appdistribution.FirebaseAppDistributionException @@ -36,15 +35,22 @@ import `in`.sminnovations.hpostesting.databinding.ActivityDashboardBinding import okhttp3.ResponseBody import java.io.File -open interface IDataCollector { +interface NatsMessageCallback { + fun onMessageReceived(topic: String, message: String) +} + +open interface IDataCollector: NatsMessageCallback { fun setConnect(connect: Boolean) fun setResponse(response: String) } + @AndroidEntryPoint class DashboardActivity : AppCompatActivity(), IDataCollector { + val TAG = "DashboardActivity" private lateinit var appBarConfiguration: AppBarConfiguration private lateinit var binding: ActivityDashboardBinding + lateinit var sharedPreferences: SharedPreferences var responses: String = "" lateinit var nats: NatsManager private var downloadId: Long = 0 @@ -57,15 +63,27 @@ class DashboardActivity : AppCompatActivity(), IDataCollector { super.attachBaseContext(newBase) } + override fun onMessageReceived(topic: String, message: String) { + // Handle incoming messages from NATS + Log.d(TAG, "Received message on topic $topic: $message") + } + override fun onCreate(savedInstanceState: Bundle?) { super.onCreate(savedInstanceState) binding = ActivityDashboardBinding.inflate(layoutInflater) + sharedPreferences = this.getSharedPreferences("HEMOCUBE", Context.MODE_PRIVATE) setContentView(binding.root) setSupportActionBar(binding.appBarDashboard.toolbar) nats = NatsManager(this) - nats.connect() - nats.pub("server.hpos.HCV-000-3001.ping", "THIS IS A TEST MSG") + if (Build.VERSION.SDK_INT >= Build.VERSION_CODES.O) { + nats.connect() + } + + val deviceId = sharedPreferences.getString(Constants.DEVICE_ID, "") + + nats.sub("server.hpos.${deviceId}.ping") + nats.pub("server.hpos.${deviceId}.ping", "THIS IS A TEST MSG") hemocubeViewModel.deviceUpdate.observe(this) { result -> when (result) {