Skip to content

Commit 2a524e3

Browse files
authored
feat: adding serialization using fory (#118)
1 parent bdf6c71 commit 2a524e3

18 files changed

Lines changed: 114 additions & 547 deletions

File tree

build.gradle.kts

Lines changed: 2 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,5 @@
1-
import com.google.protobuf.gradle.id
2-
31
plugins {
42
application
5-
id("com.google.protobuf") version "0.9.4"
63
id("com.ncorti.ktfmt.gradle") version "0.24.0"
74
id("org.pkl-lang") version "0.30.2"
85
kotlin("kapt") version "2.3.0"
@@ -39,7 +36,8 @@ dependencies {
3936
implementation("org.jgrapht:jgrapht-core:1.5.2")
4037
implementation("org.jgrapht:jgrapht-io:1.5.2")
4138

42-
implementation("com.google.protobuf:protobuf-java:4.32.0")
39+
implementation("org.apache.fory:fory-core:0.15.0")
40+
implementation("org.apache.fory:fory-kotlin:0.15.0")
4341

4442
implementation("io.etcd:jetcd-core:0.8.6")
4543
implementation("org.eclipse.zenoh:zenoh-kotlin:1.7.2")
@@ -121,14 +119,3 @@ pkl {
121119
}
122120
}
123121
}
124-
125-
protobuf {
126-
generateProtoTasks {
127-
all().forEach { task ->
128-
task.builtins {
129-
id("python")
130-
id("cpp")
131-
}
132-
}
133-
}
134-
}

compose.yaml

Lines changed: 1 addition & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -1,15 +1,4 @@
11
services:
2-
influxdb:
3-
image: "influxdb:2"
4-
ports:
5-
- "8086:8086"
6-
environment:
7-
- DOCKER_INFLUXDB_INIT_MODE=setup
8-
- DOCKER_INFLUXDB_INIT_USERNAME=admin
9-
- DOCKER_INFLUXDB_INIT_PASSWORD=adminadmin
10-
- DOCKER_INFLUXDB_INIT_ORG=org
11-
- DOCKER_INFLUXDB_INIT_BUCKET=bucket
12-
- DOCKER_INFLUXDB_INIT_ADMIN_TOKEN=bzO10KmR8x
132
etcd:
143
image: "quay.io/coreos/etcd:v3.5.7"
154
container_name: "etcd"
@@ -22,8 +11,4 @@ services:
2211
- "--listen-client-urls"
2312
- "http://0.0.0.0:2379"
2413
- "--advertise-client-urls"
25-
- "http://0.0.0.0:2379"
26-
zipkin:
27-
image: "openzipkin/zipkin:latest"
28-
ports:
29-
- "9411:9411"
14+
- "http://0.0.0.0:2379"

metrics/.gitkeep

Whitespace-only changes.

src/main/kotlin/at/ac/uibk/dps/cirrina/execution/object/EventHandler.kt

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,7 @@ package at.ac.uibk.dps.cirrina.execution.`object`
33
import at.ac.uibk.dps.cirrina.EnvironmentVariables
44
import at.ac.uibk.dps.cirrina.csm.Csml
55
import at.ac.uibk.dps.cirrina.execution.graph.EventGraph
6-
import at.ac.uibk.dps.cirrina.execution.util.EventExchange
6+
import at.ac.uibk.dps.cirrina.execution.util.Serializer
77
import io.zenoh.Config
88
import io.zenoh.Session
99
import io.zenoh.Zenoh
@@ -68,7 +68,7 @@ class EventHandler() : AutoCloseable {
6868
fun emit(event: Event) {
6969
val key = event.toKey() ?: return
7070
val publisher = publishers[key] ?: error("no publisher for topic '${key}'")
71-
val payload = ZBytes.from(EventExchange.toBytes(event))
71+
val payload = ZBytes.from(Serializer.serialize(event))
7272

7373
publisher.put(payload).onFailure { error("failed to send event '$event'") }
7474
}
@@ -104,7 +104,7 @@ class EventHandler() : AutoCloseable {
104104
private fun handleIncoming(sample: Sample) {
105105
runCatching {
106106
val bytes = sample.payload.toBytes()
107-
val event = EventExchange.fromBytes(bytes)
107+
val event = Serializer.deserialize<Event>(bytes)
108108

109109
propagate(event)
110110
}

src/main/kotlin/at/ac/uibk/dps/cirrina/execution/provider/ContextEtcd.kt

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@ package at.ac.uibk.dps.cirrina.execution.provider
22

33
import at.ac.uibk.dps.cirrina.execution.`object`.Context
44
import at.ac.uibk.dps.cirrina.execution.`object`.ContextVariable
5-
import at.ac.uibk.dps.cirrina.execution.util.ValueExchange
5+
import at.ac.uibk.dps.cirrina.execution.util.Serializer
66
import io.etcd.jetcd.ByteSequence
77
import io.etcd.jetcd.Client
88
import io.etcd.jetcd.op.Cmp
@@ -36,7 +36,7 @@ class ContextEtcd(endpoints: List<String>) : Context {
3636

3737
override fun create(name: String, value: Any?): Int {
3838
val key = name.toByteSequence()
39-
val bytes = value.toBytes()
39+
val bytes = value?.toBytes() ?: byteArrayOf()
4040

4141
val txn =
4242
client.kvClient
@@ -52,7 +52,7 @@ class ContextEtcd(endpoints: List<String>) : Context {
5252

5353
override fun assign(name: String, value: Any?): Int {
5454
val key = name.toByteSequence()
55-
val bytes = value.toBytes()
55+
val bytes = value?.toBytes() ?: byteArrayOf()
5656

5757
val txn =
5858
client.kvClient
@@ -87,9 +87,9 @@ class ContextEtcd(endpoints: List<String>) : Context {
8787

8888
private fun String.toByteSequence() = ByteSequence.from(this, StandardCharsets.UTF_8)
8989

90-
private fun Any?.toBytes(): ByteArray = ValueExchange.toBytes(this)
90+
private fun Any.toBytes(): ByteArray = Serializer.serialize(this)
9191

92-
private fun ByteArray.fromBytes(): Any? = ValueExchange.fromBytes(this)
92+
private fun ByteArray.fromBytes(): Any = Serializer.deserialize(this)
9393

9494
override fun close() {
9595
client.close()

src/main/kotlin/at/ac/uibk/dps/cirrina/execution/service/ServiceImplementation.kt

Lines changed: 3 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -4,8 +4,7 @@ import at.ac.uibk.dps.cirrina.csm.Csml.HttpMethod
44
import at.ac.uibk.dps.cirrina.csm.Csml.HttpServiceImplementationBinding
55
import at.ac.uibk.dps.cirrina.csm.Csml.ServiceImplementationBinding
66
import at.ac.uibk.dps.cirrina.execution.`object`.ContextVariable
7-
import at.ac.uibk.dps.cirrina.execution.`object`.exchange.ContextVariableProtos
8-
import at.ac.uibk.dps.cirrina.execution.util.ContextVariableExchange
7+
import at.ac.uibk.dps.cirrina.execution.util.Serializer
98
import com.google.protobuf.InvalidProtocolBufferException
109
import java.net.HttpURLConnection
1110
import java.net.URI
@@ -55,7 +54,7 @@ class HttpServiceImplementation(
5554
override suspend fun invoke(input: List<ContextVariable>): List<ContextVariable> {
5655
require(input.none { it.isLazy }) { "all variables must be evaluated before conversion" }
5756

58-
val payload = serializeInput(input)
57+
val payload = Serializer.serialize(input)
5958
val uri = URI(scheme, null, host, port, endPoint, null, null)
6059

6160
val request =
@@ -71,15 +70,6 @@ class HttpServiceImplementation(
7170
return handleResponse(response)
7271
}
7372

74-
private fun serializeInput(input: List<ContextVariable>): ByteArray {
75-
if (input.isEmpty()) return byteArrayOf()
76-
77-
return ContextVariableProtos.ContextVariables.newBuilder()
78-
.addAllData(input.map { ContextVariableExchange.toProto(it) })
79-
.build()
80-
.toByteArray()
81-
}
82-
8373
private fun handleResponse(response: HttpResponse<ByteArray>): List<ContextVariable> {
8474
val statusCode = response.statusCode()
8575

@@ -91,9 +81,7 @@ class HttpServiceImplementation(
9181
if (body == null || body.isEmpty()) return emptyList()
9282

9383
return try {
94-
ContextVariableProtos.ContextVariables.parseFrom(body).dataList.map {
95-
ContextVariableExchange.fromProto(it)
96-
}
84+
Serializer.deserialize(body)
9785
} catch (_: InvalidProtocolBufferException) {
9886
error("unexpected http service response format")
9987
}

src/main/kotlin/at/ac/uibk/dps/cirrina/execution/util/ContextVariableExchange.kt

Lines changed: 0 additions & 15 deletions
This file was deleted.

src/main/kotlin/at/ac/uibk/dps/cirrina/execution/util/EventExchange.kt

Lines changed: 0 additions & 72 deletions
This file was deleted.
Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,33 @@
1+
package at.ac.uibk.dps.cirrina.execution.util
2+
3+
import at.ac.uibk.dps.cirrina.csm.Csml
4+
import at.ac.uibk.dps.cirrina.execution.`object`.ContextVariable
5+
import at.ac.uibk.dps.cirrina.execution.`object`.Event
6+
import org.apache.fory.Fory
7+
import org.apache.fory.ThreadSafeFory
8+
import org.apache.fory.config.Language
9+
10+
object Serializer {
11+
private val fory: ThreadSafeFory =
12+
Fory.builder().withLanguage(Language.JAVA).buildThreadSafeFory().apply {
13+
register(Event::class.java)
14+
register(Csml.EventChannel::class.java)
15+
register(ContextVariable::class.java)
16+
}
17+
18+
fun serialize(obj: Any): ByteArray {
19+
if (obj is Event) {
20+
val data = obj.data
21+
for (i in 0 until data.size) {
22+
if (data[i].isLazy) error("event '${obj.topic}' has unevaluated data")
23+
}
24+
}
25+
26+
return fory.serialize(obj)
27+
}
28+
29+
@Suppress("UNCHECKED_CAST")
30+
fun <T> deserialize(data: ByteArray): T {
31+
return fory.deserialize(data) as T
32+
}
33+
}

src/main/kotlin/at/ac/uibk/dps/cirrina/execution/util/ValueExchange.kt

Lines changed: 0 additions & 88 deletions
This file was deleted.

0 commit comments

Comments
 (0)