handle method and dynamic config for NATS

This commit is contained in:
Pritimay Sarkar
2024-02-13 08:59:03 +05:30
parent 6577842ba9
commit 23cf38445f
5 changed files with 139 additions and 120 deletions

View File

@@ -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

View File

@@ -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 = ""
)

View File

@@ -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")

View File

@@ -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)
}
}

View File

@@ -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) {