summaryrefslogtreecommitdiff
path: root/simulator/opendc
diff options
context:
space:
mode:
authorFabian Mastenbroek <mail.fabianm@gmail.com>2020-10-01 00:23:37 +0200
committerFabian Mastenbroek <mail.fabianm@gmail.com>2020-10-01 00:23:37 +0200
commit0df646c2951e9950f27472fdf0cb2624303c2d74 (patch)
tree0e8fed649af780a518afb8e60a7823d212da213f /simulator/opendc
parentfcae560208df4860bc7461f955bf3b522b0e61c5 (diff)
Move custom Flows to separate utility module
This change moves the custom Flow object we provide (e.g. EventFlow and StateFlow) outside of the odcsim-api module into a separate opendc-utils module. This is in preparation for the removal of the odcsim components in OpenDC.
Diffstat (limited to 'simulator/opendc')
-rw-r--r--simulator/opendc/opendc-compute/build.gradle.kts1
-rw-r--r--simulator/opendc/opendc-compute/src/main/kotlin/com/atlarge/opendc/compute/metal/driver/SimpleBareMetalDriver.kt4
-rw-r--r--simulator/opendc/opendc-compute/src/main/kotlin/com/atlarge/opendc/compute/virt/driver/SimpleVirtDriver.kt2
-rw-r--r--simulator/opendc/opendc-compute/src/main/kotlin/com/atlarge/opendc/compute/virt/service/SimpleVirtProvisioningService.kt2
-rw-r--r--simulator/opendc/opendc-utils/build.gradle.kts32
-rw-r--r--simulator/opendc/opendc-utils/src/main/kotlin/org/opendc/utils/flow/EventFlow.kt96
-rw-r--r--simulator/opendc/opendc-utils/src/main/kotlin/org/opendc/utils/flow/StateFlow.kt79
-rw-r--r--simulator/opendc/opendc-workflows/build.gradle.kts1
-rw-r--r--simulator/opendc/opendc-workflows/src/main/kotlin/com/atlarge/opendc/workflows/service/StageWorkflowService.kt2
9 files changed, 214 insertions, 5 deletions
diff --git a/simulator/opendc/opendc-compute/build.gradle.kts b/simulator/opendc/opendc-compute/build.gradle.kts
index 0e44785e..a8852cba 100644
--- a/simulator/opendc/opendc-compute/build.gradle.kts
+++ b/simulator/opendc/opendc-compute/build.gradle.kts
@@ -32,6 +32,7 @@ plugins {
dependencies {
api(project(":odcsim:odcsim-api"))
api(project(":opendc:opendc-core"))
+ implementation(project(":opendc:opendc-utils"))
implementation("io.github.microutils:kotlin-logging:1.7.9")
testImplementation(project(":opendc:opendc-simulator"))
diff --git a/simulator/opendc/opendc-compute/src/main/kotlin/com/atlarge/opendc/compute/metal/driver/SimpleBareMetalDriver.kt b/simulator/opendc/opendc-compute/src/main/kotlin/com/atlarge/opendc/compute/metal/driver/SimpleBareMetalDriver.kt
index c118cc3d..0c758e6b 100644
--- a/simulator/opendc/opendc-compute/src/main/kotlin/com/atlarge/opendc/compute/metal/driver/SimpleBareMetalDriver.kt
+++ b/simulator/opendc/opendc-compute/src/main/kotlin/com/atlarge/opendc/compute/metal/driver/SimpleBareMetalDriver.kt
@@ -24,8 +24,6 @@
package com.atlarge.opendc.compute.metal.driver
-import com.atlarge.odcsim.flow.EventFlow
-import com.atlarge.odcsim.flow.StateFlow
import com.atlarge.opendc.compute.core.Flavor
import com.atlarge.opendc.compute.core.MemoryUnit
import com.atlarge.opendc.compute.core.ProcessingUnit
@@ -57,6 +55,8 @@ import kotlinx.coroutines.intrinsics.startCoroutineCancellable
import kotlinx.coroutines.launch
import kotlinx.coroutines.selects.SelectClause0
import kotlinx.coroutines.selects.SelectInstance
+import org.opendc.utils.flow.EventFlow
+import org.opendc.utils.flow.StateFlow
import java.lang.Exception
import java.time.Clock
import java.util.UUID
diff --git a/simulator/opendc/opendc-compute/src/main/kotlin/com/atlarge/opendc/compute/virt/driver/SimpleVirtDriver.kt b/simulator/opendc/opendc-compute/src/main/kotlin/com/atlarge/opendc/compute/virt/driver/SimpleVirtDriver.kt
index bd266208..fa172e6e 100644
--- a/simulator/opendc/opendc-compute/src/main/kotlin/com/atlarge/opendc/compute/virt/driver/SimpleVirtDriver.kt
+++ b/simulator/opendc/opendc-compute/src/main/kotlin/com/atlarge/opendc/compute/virt/driver/SimpleVirtDriver.kt
@@ -24,7 +24,6 @@
package com.atlarge.opendc.compute.virt.driver
-import com.atlarge.odcsim.flow.EventFlow
import com.atlarge.opendc.compute.core.Flavor
import com.atlarge.opendc.compute.core.ProcessingUnit
import com.atlarge.opendc.compute.core.Server
@@ -47,6 +46,7 @@ import kotlinx.coroutines.selects.SelectClause0
import kotlinx.coroutines.selects.SelectInstance
import kotlinx.coroutines.selects.select
import mu.KotlinLogging
+import org.opendc.utils.flow.EventFlow
import java.time.Clock
import java.util.UUID
import kotlin.math.ceil
diff --git a/simulator/opendc/opendc-compute/src/main/kotlin/com/atlarge/opendc/compute/virt/service/SimpleVirtProvisioningService.kt b/simulator/opendc/opendc-compute/src/main/kotlin/com/atlarge/opendc/compute/virt/service/SimpleVirtProvisioningService.kt
index 6b2cfc40..5e151cbb 100644
--- a/simulator/opendc/opendc-compute/src/main/kotlin/com/atlarge/opendc/compute/virt/service/SimpleVirtProvisioningService.kt
+++ b/simulator/opendc/opendc-compute/src/main/kotlin/com/atlarge/opendc/compute/virt/service/SimpleVirtProvisioningService.kt
@@ -1,6 +1,5 @@
package com.atlarge.opendc.compute.virt.service
-import com.atlarge.odcsim.flow.EventFlow
import com.atlarge.opendc.compute.core.Flavor
import com.atlarge.opendc.compute.core.Server
import com.atlarge.opendc.compute.core.ServerEvent
@@ -19,6 +18,7 @@ import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.launchIn
import kotlinx.coroutines.flow.onEach
import mu.KotlinLogging
+import org.opendc.utils.flow.EventFlow
import java.time.Clock
import kotlin.coroutines.Continuation
import kotlin.coroutines.resume
diff --git a/simulator/opendc/opendc-utils/build.gradle.kts b/simulator/opendc/opendc-utils/build.gradle.kts
new file mode 100644
index 00000000..d66148c4
--- /dev/null
+++ b/simulator/opendc/opendc-utils/build.gradle.kts
@@ -0,0 +1,32 @@
+/*
+ * Copyright (c) 2020 AtLarge Research
+ *
+ * Permission is hereby granted, free of charge, to any person obtaining a copy
+ * of this software and associated documentation files (the "Software"), to deal
+ * in the Software without restriction, including without limitation the rights
+ * to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
+ * copies of the Software, and to permit persons to whom the Software is
+ * furnished to do so, subject to the following conditions:
+ *
+ * The above copyright notice and this permission notice shall be included in all
+ * copies or substantial portions of the Software.
+ *
+ * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
+ * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
+ * FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
+ * AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
+ * LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
+ * OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
+ * SOFTWARE.
+ */
+
+description = "Utilities used across OpenDC modules"
+
+/* Build configuration */
+plugins {
+ `kotlin-library-convention`
+}
+
+dependencies {
+ api("org.jetbrains.kotlinx:kotlinx-coroutines-core:${Library.KOTLINX_COROUTINES}")
+}
diff --git a/simulator/opendc/opendc-utils/src/main/kotlin/org/opendc/utils/flow/EventFlow.kt b/simulator/opendc/opendc-utils/src/main/kotlin/org/opendc/utils/flow/EventFlow.kt
new file mode 100644
index 00000000..948595b1
--- /dev/null
+++ b/simulator/opendc/opendc-utils/src/main/kotlin/org/opendc/utils/flow/EventFlow.kt
@@ -0,0 +1,96 @@
+/*
+ * Copyright (c) 2020 AtLarge Research
+ *
+ * Permission is hereby granted, free of charge, to any person obtaining a copy
+ * of this software and associated documentation files (the "Software"), to deal
+ * in the Software without restriction, including without limitation the rights
+ * to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
+ * copies of the Software, and to permit persons to whom the Software is
+ * furnished to do so, subject to the following conditions:
+ *
+ * The above copyright notice and this permission notice shall be included in all
+ * copies or substantial portions of the Software.
+ *
+ * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
+ * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
+ * FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
+ * AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
+ * LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
+ * OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
+ * SOFTWARE.
+ */
+
+package org.opendc.utils.flow
+
+import kotlinx.coroutines.ExperimentalCoroutinesApi
+import kotlinx.coroutines.FlowPreview
+import kotlinx.coroutines.InternalCoroutinesApi
+import kotlinx.coroutines.channels.Channel
+import kotlinx.coroutines.channels.SendChannel
+import kotlinx.coroutines.flow.Flow
+import kotlinx.coroutines.flow.FlowCollector
+import kotlinx.coroutines.flow.consumeAsFlow
+
+/**
+ * A [Flow] that can be used to emit events.
+ */
+public interface EventFlow<T> : Flow<T> {
+ /**
+ * Emit the specified [event].
+ */
+ public fun emit(event: T)
+
+ /**
+ * Close the flow.
+ */
+ public fun close()
+}
+
+/**
+ * Creates a new [EventFlow].
+ */
+@Suppress("FunctionName")
+public fun <T> EventFlow(): EventFlow<T> = EventFlowImpl()
+
+/**
+ * Internal implementation of the [EventFlow] class.
+ */
+@OptIn(ExperimentalCoroutinesApi::class, FlowPreview::class)
+private class EventFlowImpl<T> : EventFlow<T> {
+ private var closed: Boolean = false
+ private val subscribers = HashMap<SendChannel<T>, Unit>()
+
+ override fun emit(event: T) {
+ synchronized(this) {
+ for ((chan, _) in subscribers) {
+ chan.offer(event)
+ }
+ }
+ }
+
+ override fun close() {
+ synchronized(this) {
+ closed = true
+
+ for ((chan, _) in subscribers) {
+ chan.close()
+ }
+ }
+ }
+
+ @InternalCoroutinesApi
+ override suspend fun collect(collector: FlowCollector<T>) {
+ val channel: Channel<T>
+ synchronized(this) {
+ if (closed) {
+ return
+ }
+
+ channel = Channel(Channel.UNLIMITED)
+ subscribers[channel] = Unit
+ }
+ channel.consumeAsFlow().collect(collector)
+ }
+
+ override fun toString(): String = "EventFlow"
+}
diff --git a/simulator/opendc/opendc-utils/src/main/kotlin/org/opendc/utils/flow/StateFlow.kt b/simulator/opendc/opendc-utils/src/main/kotlin/org/opendc/utils/flow/StateFlow.kt
new file mode 100644
index 00000000..996e7700
--- /dev/null
+++ b/simulator/opendc/opendc-utils/src/main/kotlin/org/opendc/utils/flow/StateFlow.kt
@@ -0,0 +1,79 @@
+/*
+ * Copyright (c) 2020 AtLarge Research
+ *
+ * Permission is hereby granted, free of charge, to any person obtaining a copy
+ * of this software and associated documentation files (the "Software"), to deal
+ * in the Software without restriction, including without limitation the rights
+ * to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
+ * copies of the Software, and to permit persons to whom the Software is
+ * furnished to do so, subject to the following conditions:
+ *
+ * The above copyright notice and this permission notice shall be included in all
+ * copies or substantial portions of the Software.
+ *
+ * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
+ * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
+ * FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
+ * AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
+ * LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
+ * OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
+ * SOFTWARE.
+ */
+
+package org.opendc.utils.flow
+
+import kotlinx.coroutines.ExperimentalCoroutinesApi
+import kotlinx.coroutines.FlowPreview
+import kotlinx.coroutines.InternalCoroutinesApi
+import kotlinx.coroutines.channels.BroadcastChannel
+import kotlinx.coroutines.channels.ConflatedBroadcastChannel
+import kotlinx.coroutines.flow.Flow
+import kotlinx.coroutines.flow.FlowCollector
+import kotlinx.coroutines.flow.asFlow
+
+/**
+ * A [Flow] that contains a single value that changes over time.
+ *
+ * This class exists to implement the DataFlow/StateFlow functionality that will be implemented in `kotlinx-coroutines`
+ * in the future, but is not available yet.
+ * See: https://github.com/Kotlin/kotlinx.coroutines/pull/1354
+ */
+public interface StateFlow<T> : Flow<T> {
+ /**
+ * The current value of this flow.
+ *
+ * Setting a value that is [equal][Any.equals] to the previous one does nothing.
+ */
+ public var value: T
+}
+
+/**
+ * Creates a [StateFlow] with a given initial [value].
+ */
+@Suppress("FunctionName")
+public fun <T> StateFlow(value: T): StateFlow<T> = StateFlowImpl(value)
+
+/**
+ * Internal implementation of the [StateFlow] interface.
+ */
+@OptIn(ExperimentalCoroutinesApi::class, FlowPreview::class)
+private class StateFlowImpl<T>(initialValue: T) : StateFlow<T> {
+ /**
+ * The [BroadcastChannel] to back this flow.
+ */
+ private val chan = ConflatedBroadcastChannel(initialValue)
+
+ /**
+ * The internal [Flow] backing this flow.
+ */
+ private val flow = chan.asFlow()
+
+ public override var value: T = initialValue
+ set(value) {
+ chan.offer(value)
+ field = value
+ }
+
+ @InternalCoroutinesApi
+ override suspend fun collect(collector: FlowCollector<T>) = flow.collect(collector)
+}
diff --git a/simulator/opendc/opendc-workflows/build.gradle.kts b/simulator/opendc/opendc-workflows/build.gradle.kts
index 62c4bc25..f8a9a1f3 100644
--- a/simulator/opendc/opendc-workflows/build.gradle.kts
+++ b/simulator/opendc/opendc-workflows/build.gradle.kts
@@ -32,6 +32,7 @@ plugins {
dependencies {
api(project(":opendc:opendc-core"))
api(project(":opendc:opendc-compute"))
+ implementation(project(":opendc:opendc-utils"))
testImplementation(project(":opendc:opendc-simulator"))
testImplementation(project(":opendc:opendc-format"))
diff --git a/simulator/opendc/opendc-workflows/src/main/kotlin/com/atlarge/opendc/workflows/service/StageWorkflowService.kt b/simulator/opendc/opendc-workflows/src/main/kotlin/com/atlarge/opendc/workflows/service/StageWorkflowService.kt
index aea27972..3a5b963c 100644
--- a/simulator/opendc/opendc-workflows/src/main/kotlin/com/atlarge/opendc/workflows/service/StageWorkflowService.kt
+++ b/simulator/opendc/opendc-workflows/src/main/kotlin/com/atlarge/opendc/workflows/service/StageWorkflowService.kt
@@ -24,7 +24,6 @@
package com.atlarge.opendc.workflows.service
-import com.atlarge.odcsim.flow.EventFlow
import com.atlarge.opendc.compute.core.Server
import com.atlarge.opendc.compute.core.ServerEvent
import com.atlarge.opendc.compute.core.ServerState
@@ -43,6 +42,7 @@ import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.launchIn
import kotlinx.coroutines.flow.onEach
import kotlinx.coroutines.launch
+import org.opendc.utils.flow.EventFlow
import java.time.Clock
import java.util.PriorityQueue
import java.util.Queue