[ref][fix] refactored Publisher trait

- now uses watchFrame instead of vector
 - all BroadcastObject uses are moved to interfaces
 - [fix] fixed missing writeLock in BroadcastObject
This commit is contained in:
2020-01-29 01:19:11 -06:00
parent 56b46d4dfb
commit 75988aa521
5 changed files with 27 additions and 118 deletions

View File

@@ -1,7 +1,6 @@
package scim.components package scim.components
import scim.datastruct.{Behaviour, BehaviourVector, Cell, Configuration, FunctionArgs, GroundTruth, Group, Sensor, SensorCount, SensorResult, SharedMemory, SimulationBehaviours, SimulationConfig, Time, Truth} import scim.datastruct.{Behaviour, BehaviourVector, Cell, Configuration, FunctionArgs, GroundTruth, Group, Sensor, SensorCount, SensorResult, SharedMemory, SimulationBehaviours, SimulationConfig, Time, Truth}
import scim.lib.ConcurrentBroadcast.BroadcastObject
import scim.lib.debugUtil.{DebugType, debugf} import scim.lib.debugUtil.{DebugType, debugf}
import scim.lib.mapUtil._ import scim.lib.mapUtil._
import scim.lib.simInterfaces.SimulationBehaviour import scim.lib.simInterfaces.SimulationBehaviour
@@ -68,7 +67,7 @@ class SimulationConfigParser(input: Map[String, Any], simulationBehaviours: Simu
val sharedMemoryTruthMap = traversableToMap[Time, Map[Cell, Map[Group, Truth]]](sharedMemoryTruth) val sharedMemoryTruthMap = traversableToMap[Time, Map[Cell, Map[Group, Truth]]](sharedMemoryTruth)
val sharedMemory = SharedMemory(sharedMemorySensorMap, sharedMemoryTruthMap, groundTruthValues.map(kv => kv._1 -> kv._2.toMap).toMap, config, groups.toMap, Vector.tabulate(config.duration)(_ => new BroadcastObject[Boolean]())) val sharedMemory = SharedMemory(sharedMemorySensorMap, sharedMemoryTruthMap, groundTruthValues.map(kv => kv._1 -> kv._2.toMap).toMap, config, groups.toMap)
Configuration(configurationMapping, sharedMemory) Configuration(configurationMapping, sharedMemory)
} }

View File

@@ -46,7 +46,7 @@ class Simulation(config: Configuration) {
cells.foreach(simCell) cells.foreach(simCell)
} }
// broadcast the results // broadcast the results
sharedMemory.publish(time) sharedMemory.publish()
} }
private def simGroups(time: Time, cell: Cell, groups: Map[Group, List[(Group, Sensor, Option[GroundTruth])]], configTruths: Option[Map[Group, Truth]], prevTruths: Map[Group, Truth], resultTruths: Map[Group, Truth], resultSensors: Map[Group, Vector[SensorResult]]): Unit = { private def simGroups(time: Time, cell: Cell, groups: Map[Group, List[(Group, Sensor, Option[GroundTruth])]], configTruths: Option[Map[Group, Truth]], prevTruths: Map[Group, Truth], resultTruths: Map[Group, Truth], resultSensors: Map[Group, Vector[SensorResult]]): Unit = {

View File

@@ -1,6 +1,5 @@
package scim.datastruct package scim.datastruct
import scim.lib.ConcurrentBroadcast._
import scim.lib.simInterfaces._ import scim.lib.simInterfaces._
import scim.lib.configInterfaces._ import scim.lib.configInterfaces._
import scim.lib.mapUtil._ import scim.lib.mapUtil._
@@ -31,8 +30,7 @@ case class SharedMemory(sensor: Map[Time, Map[Cell, Map[Group, Vector[SensorResu
groundTruth: Map[Time, Map[Cell, Map[Group, Truth]]], groundTruth: Map[Time, Map[Cell, Map[Group, Truth]]],
groundTruthConfig: Map[Time, Map[Group, Truth]], groundTruthConfig: Map[Time, Map[Group, Truth]],
config: SimulationConfig, config: SimulationConfig,
groups: Map[String, Group], groups: Map[String, Group]) extends Publisher {
private val output: Vector[BroadcastObject[Boolean]]) extends Publisher {
// memory optimization: remove sensor and throw output into output // memory optimization: remove sensor and throw output into output
// only map Truths from config into groundtruth, do a if groundTruth.contains(time) groundTruth(time) else prevTruth in case of time.start Truth(None) (or better map it to time.start anyways in parser) // only map Truths from config into groundtruth, do a if groundTruth.contains(time) groundTruth(time) else prevTruth in case of time.start Truth(None) (or better map it to time.start anyways in parser)
// can drop cell from ground truth // can drop cell from ground truth
@@ -55,19 +53,17 @@ case class SharedMemory(sensor: Map[Time, Map[Cell, Map[Group, Vector[SensorResu
def getResult(time: Time): Map[Cell, Map[Group, Vector[SensorResult]]] = { def getResult(time: Time): Map[Cell, Map[Group, Vector[SensorResult]]] = {
sensor(time) sensor(time)
} }
def outputReset(): Unit = { def outputReset(): Unit = {
sensor.foreach(kv => kv._2.foreach(kv2 => kv2._2.foreach(kv3 => kv3._2.foreach(res => res.reset())))) sensor.foreach(kv => kv._2.foreach(kv2 => kv2._2.foreach(kv3 => kv3._2.foreach(res => res.reset()))))
groundTruth.foreach(kv => kv._2.foreach(kv2 => kv2._2.foreach(kv3 => kv3._2.reset()))) groundTruth.foreach(kv => kv._2.foreach(kv2 => kv2._2.foreach(kv3 => kv3._2.reset())))
output.foreach(e => e.broadcastReset()) reset()
} }
def subscribe(time: Time): SimulationOutput = { def subscribe(time: Time): SimulationOutput = {
output(time.time).watch() waitForOutput(time)
(getResult(time), getTruth(time)) (getResult(time), getTruth(time))
} }
def publish(time: Time): Unit = {
output(time.time).broadcast(true)
}
} }
case class SimulationBehaviours(functions: Map[String, SimulationBehaviour]) { case class SimulationBehaviours(functions: Map[String, SimulationBehaviour]) {

View File

@@ -83,7 +83,7 @@ object ConcurrentBroadcast {
private var broadcastEnded = false private var broadcastEnded = false
private var value: Option[A] = default private var value: Option[A] = default
// NON-BLOCKING: check if value has been set, write new value and return if it was overwritten // NON-BLOCKING: write new value
def broadcast(newValue: A): Unit = { def broadcast(newValue: A): Unit = {
writeLock(() => { writeLock(() => {
value = Some(newValue) value = Some(newValue)
@@ -118,9 +118,11 @@ object ConcurrentBroadcast {
} }
def broadcastReset(): Unit = { def broadcastReset(): Unit = {
broadcastEnded = false writeLock(() => {
frame = 0 broadcastEnded = false
value = default frame = 0
value = default
})
} }
// BLOCKING: watch value, block until it is set, then return it // BLOCKING: watch value, block until it is set, then return it
@@ -174,101 +176,3 @@ object ConcurrentBroadcast {
} }
} }
class Test2 {
import scim.lib.ConcurrentBroadcast._
private val stream: mutable.ArraySeq[BroadcastObject[Int]] = new mutable.ArraySeq[Int](5).map(_ => new BroadcastObject[Int])
class Producer(stream: mutable.ArraySeq[BroadcastObject[Int]]) extends Broadcast[Int](stream) with Broadcaster[Int] {
def a(): Unit = {
stream.foreach(element => {element.broadcast(1); Thread.sleep(1000)})
}
override def broadcast(stream: mutable.ArraySeq[BroadcastObject[Int]]): Unit = {
val f: (Int, Int) => Int = (x, y) => x + y
var previous: Int = 0
stream.foreach(element => {
previous = f(previous, 2)
element.broadcast(previous)
})
}
}
class Watcher(id: Int, stream: mutable.ArraySeq[BroadcastObject[Int]]) extends Runnable {
def run(): Unit = {
stream.foreach(element => { Thread.sleep(id); println(id + ": " + element.watch() + " - " + Calendar.getInstance().getTime()) })
}
}
def run(): Unit = {
val threadPool = Executors.newFixedThreadPool(4)
val producer = new Producer(stream)
val watcher = new Watcher(0,stream)
val watcher2 = new Watcher(3000,stream)
val p = new Thread(producer)
val t = new Thread(watcher)
val t2 = new Thread(watcher2)
println("start producer")
p.start()
println("start watcher")
t.start()
t2.start()
println("join producer")
p.join()
println("join watcher")
t.join()
t2.join()
println("done")
}
}
class Test {
import scim.lib.ConcurrentBroadcast.BroadcastObject
private val stream: mutable.ArraySeq[BroadcastObject[Int]] = new mutable.ArraySeq[Int](5).map(_ => new BroadcastObject[Int])
class Producer(stream: mutable.ArraySeq[BroadcastObject[Int]]) extends Runnable {
def run(): Unit = {
stream.foreach(element => {element.broadcast(1); Thread.sleep(1000)})
}
}
class Watcher(id: Int, stream: mutable.ArraySeq[BroadcastObject[Int]]) extends Runnable {
def run(): Unit = {
stream.foreach(element => { Thread.sleep(id); println(id + ": " + element.watch() + " - " + Calendar.getInstance().getTime()) })
}
}
def run(): Unit = {
val threadPool = Executors.newFixedThreadPool(4)
val producer = new Producer(stream)
val watcher = new Watcher(0,stream)
val watcher2 = new Watcher(3000,stream)
val p = new Thread(producer)
val t = new Thread(watcher)
val t2 = new Thread(watcher2)
println("start producer")
p.start()
println("start watcher")
t.start()
t2.start()
println("join producer")
p.join()
println("join watcher")
t.join()
t2.join()
println("done")
}
}

View File

@@ -1,23 +1,33 @@
package scim.lib package scim.lib
import scim.datastruct.{Behaviour, Cell, Group, SensorCount, SensorResult, Time, Truth} import scim.datastruct.{Behaviour, Cell, FunctionArgs, Group, SensorCount, SensorResult, Time, Truth}
import scim.lib.simInterfaces.SimulationOutput import scim.lib.simInterfaces.SimulationOutput
import scala.collection.parallel.ForkJoinTaskSupport import scala.collection.parallel.ForkJoinTaskSupport
package object simInterfaces { package object simInterfaces {
import scim.datastruct._
// maybe change to by-name parameters? have to check performance gain // maybe change to by-name parameters? have to check performance gain
type SimulationBehaviour = (Cell, Time, Option[Double], Option[FunctionArgs]) => Option[Double] type SimulationBehaviour = (Cell, Time, Option[Double], Option[FunctionArgs]) => Option[Double]
type SimulationOutput = (Map[Cell, Map[Group, Vector[SensorResult]]], Map[Cell, Map[Group, Truth]]) type SimulationOutput = (Map[Cell, Map[Group, Vector[SensorResult]]], Map[Cell, Map[Group, Truth]])
} }
package object configInterfaces { package object configInterfaces {
// have to import here, get compiler error otherwise
import scim.lib.ConcurrentBroadcast.BroadcastObject
trait Publisher { trait Publisher {
private val output = new BroadcastObject[Boolean]()
protected def waitForOutput(time: Time): Unit = {
// frames start at 1, simulation ticks at 0
output.watchFrame(time.time + 1)
}
def subscribe(time: Time): SimulationOutput def subscribe(time: Time): SimulationOutput
def publish(time: Time): Unit def publish(): Unit = {
def outputReset(): Unit output.broadcast(true)
}
def reset(): Unit = {
output.broadcastReset()
}
} }
trait SimConfig { trait SimConfig {
@@ -30,7 +40,7 @@ package object configInterfaces {
} }
trait Result[A] { trait Result[A] {
val result = new scim.lib.ConcurrentBroadcast.BroadcastObject[A]() val result = new BroadcastObject[A]()
def get: A = { def get: A = {
result.watch() result.watch()
} }