Loading src/dlt/gateway/src/main/kotlin/fabric/FabricConnector.kt +36 −20 Original line number Diff line number Diff line Loading @@ -37,7 +37,9 @@ package fabric import dlt.DltGateway import context.ContextOuterClass import dlt.DltGateway.DltRecord import dlt.DltGateway.DltRecordEvent import kotlinx.coroutines.channels.Channel import kotlinx.coroutines.runBlocking import org.hyperledger.fabric.gateway.Contract Loading @@ -56,7 +58,7 @@ class FabricConnector(val config: Config.DltConfig) { private val wallet: Wallet private val contract: Contract private val channels: MutableList<Channel<DltGateway.DltRecordEvent>> = mutableListOf() private val channels: MutableList<Channel<DltRecordEvent>> = mutableListOf() init { // Create a CA client for interacting with the CA. Loading @@ -78,11 +80,29 @@ class FabricConnector(val config: Config.DltConfig) { val consumer = Consumer { event: ContractEvent? -> run { println("new event detected") val parsedEvent = DltGateway.DltRecordEvent.parseFrom(event?.payload?.get()) println(parsedEvent.recordId.recordUuid) val record = DltRecord.parseFrom(event?.payload?.get()) println(record.recordId.recordUuid) val eventType: ContextOuterClass.EventTypeEnum = when (event?.name) { "Add" -> ContextOuterClass.EventTypeEnum.EVENTTYPE_CREATE "Update" -> ContextOuterClass.EventTypeEnum.EVENTTYPE_UPDATE "Remove" -> ContextOuterClass.EventTypeEnum.EVENTTYPE_REMOVE else -> ContextOuterClass.EventTypeEnum.EVENTTYPE_UNDEFINED } val pbEvent = DltRecordEvent.newBuilder() .setEvent( ContextOuterClass.Event.newBuilder() .setTimestamp( ContextOuterClass.Timestamp.newBuilder() .setTimestamp(System.currentTimeMillis().toDouble()) ) .setEventType(eventType) ) .setRecordId(record.recordId) .build() runBlocking { channels.forEach { it.trySend(parsedEvent) it.trySend(pbEvent) } } } Loading @@ -96,49 +116,45 @@ class FabricConnector(val config: Config.DltConfig) { return getContract(config, wallet) } fun putData(record: DltGateway.DltRecord): String { fun putData(record: DltRecord): String { println(record.toString()) println("Put: ${record.toByteArray().decodeToString().length}") return String( contract.submitTransaction( "AddRecord", record.recordId.domainUuid.uuid, record.recordId.recordUuid.uuid, record.recordId.type.number.toString(), record.dataJson record.toByteArray().decodeToString() ) ) } fun getData(uuid: String): DltGateway.DltRecord { fun getData(uuid: String): DltRecord { val result = contract.evaluateTransaction("GetRecord", uuid) return DltGateway.DltRecord.parseFrom(result) println("Get: ${result.size}") return DltRecord.parseFrom(result) } fun updateData(record: DltGateway.DltRecord): String { fun updateData(record: DltRecord): String { return String( contract.submitTransaction( "UpdateRecord", record.recordId.domainUuid.uuid, record.recordId.recordUuid.uuid, record.recordId.type.number.toString(), record.dataJson record.toByteArray().decodeToString() ) ) } fun deleteData(record: DltGateway.DltRecord): String { fun deleteData(record: DltRecord): String { return String( contract.submitTransaction( "DeleteRecord", record.recordId.domainUuid.uuid, record.recordId.recordUuid.uuid, record.recordId.type.number.toString() ) ) } fun subscribeForEvents(): Channel<DltGateway.DltRecordEvent> { val produceCh = Channel<DltGateway.DltRecordEvent>() fun subscribeForEvents(): Channel<DltRecordEvent> { val produceCh = Channel<DltRecordEvent>() channels.add(produceCh) return produceCh } Loading Loading
src/dlt/gateway/src/main/kotlin/fabric/FabricConnector.kt +36 −20 Original line number Diff line number Diff line Loading @@ -37,7 +37,9 @@ package fabric import dlt.DltGateway import context.ContextOuterClass import dlt.DltGateway.DltRecord import dlt.DltGateway.DltRecordEvent import kotlinx.coroutines.channels.Channel import kotlinx.coroutines.runBlocking import org.hyperledger.fabric.gateway.Contract Loading @@ -56,7 +58,7 @@ class FabricConnector(val config: Config.DltConfig) { private val wallet: Wallet private val contract: Contract private val channels: MutableList<Channel<DltGateway.DltRecordEvent>> = mutableListOf() private val channels: MutableList<Channel<DltRecordEvent>> = mutableListOf() init { // Create a CA client for interacting with the CA. Loading @@ -78,11 +80,29 @@ class FabricConnector(val config: Config.DltConfig) { val consumer = Consumer { event: ContractEvent? -> run { println("new event detected") val parsedEvent = DltGateway.DltRecordEvent.parseFrom(event?.payload?.get()) println(parsedEvent.recordId.recordUuid) val record = DltRecord.parseFrom(event?.payload?.get()) println(record.recordId.recordUuid) val eventType: ContextOuterClass.EventTypeEnum = when (event?.name) { "Add" -> ContextOuterClass.EventTypeEnum.EVENTTYPE_CREATE "Update" -> ContextOuterClass.EventTypeEnum.EVENTTYPE_UPDATE "Remove" -> ContextOuterClass.EventTypeEnum.EVENTTYPE_REMOVE else -> ContextOuterClass.EventTypeEnum.EVENTTYPE_UNDEFINED } val pbEvent = DltRecordEvent.newBuilder() .setEvent( ContextOuterClass.Event.newBuilder() .setTimestamp( ContextOuterClass.Timestamp.newBuilder() .setTimestamp(System.currentTimeMillis().toDouble()) ) .setEventType(eventType) ) .setRecordId(record.recordId) .build() runBlocking { channels.forEach { it.trySend(parsedEvent) it.trySend(pbEvent) } } } Loading @@ -96,49 +116,45 @@ class FabricConnector(val config: Config.DltConfig) { return getContract(config, wallet) } fun putData(record: DltGateway.DltRecord): String { fun putData(record: DltRecord): String { println(record.toString()) println("Put: ${record.toByteArray().decodeToString().length}") return String( contract.submitTransaction( "AddRecord", record.recordId.domainUuid.uuid, record.recordId.recordUuid.uuid, record.recordId.type.number.toString(), record.dataJson record.toByteArray().decodeToString() ) ) } fun getData(uuid: String): DltGateway.DltRecord { fun getData(uuid: String): DltRecord { val result = contract.evaluateTransaction("GetRecord", uuid) return DltGateway.DltRecord.parseFrom(result) println("Get: ${result.size}") return DltRecord.parseFrom(result) } fun updateData(record: DltGateway.DltRecord): String { fun updateData(record: DltRecord): String { return String( contract.submitTransaction( "UpdateRecord", record.recordId.domainUuid.uuid, record.recordId.recordUuid.uuid, record.recordId.type.number.toString(), record.dataJson record.toByteArray().decodeToString() ) ) } fun deleteData(record: DltGateway.DltRecord): String { fun deleteData(record: DltRecord): String { return String( contract.submitTransaction( "DeleteRecord", record.recordId.domainUuid.uuid, record.recordId.recordUuid.uuid, record.recordId.type.number.toString() ) ) } fun subscribeForEvents(): Channel<DltGateway.DltRecordEvent> { val produceCh = Channel<DltGateway.DltRecordEvent>() fun subscribeForEvents(): Channel<DltRecordEvent> { val produceCh = Channel<DltRecordEvent>() channels.add(produceCh) return produceCh } Loading