/* * Copyright (C) 2025 The Android Open Source Project * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. * You may obtain a copy of the License at * * http://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, software * distributed under the License is distributed on an "AS IS" BASIS, * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. * See the License for the specific language governing permissions and * limitations under the License. */ package com.android.app.concurrent.benchmark.event import com.android.app.concurrent.benchmark.util.ThreadFactory import com.android.app.concurrent.benchmark.util.dbg import com.android.app.concurrent.benchmark.util.instanceName import java.util.concurrent.CancellationException import java.util.concurrent.Executor import java.util.concurrent.atomic.AtomicReference import kotlin.coroutines.Continuation import kotlin.coroutines.resume import kotlin.coroutines.resumeWithException import kotlin.coroutines.suspendCoroutine import org.junit.Assert.fail typealias SimpleEventBoxIn = EventBox> typealias SimpleObservableBoxOut = EventBox> typealias SimpleObservableStateBoxOut = EventBox> private class SimpleEventBoxImpl>(event: R) : AbstractEventBox(event) private fun SimpleEventBoxIn.unbox(): SimpleEvent = (this as SimpleEventBoxImpl).event fun interface EventListener { fun notify(value: T) } interface SimpleEvent { fun listen(listener: EventListener): AutoCloseable } interface SimplePublisher { fun publish(value: T) } open class SimpleEventImpl : SimpleEvent { private val listeners = mutableListOf>() @Synchronized override fun listen(listener: EventListener): AutoCloseable { dbg { "listeners.add(${listener.instanceName()})" } listeners.add(listener) return AutoCloseable { synchronized(listeners) { dbg { "listeners.remove(${listener.instanceName()})" } listeners.remove(listener) } } } @Synchronized protected fun notifyAll(value: T) { dbg { "notifyAll($value) {{{" } listeners.forEach { listener -> dbg { "listener.notify(${listener.instanceName()}) -> $value" } listener.notify(value) } dbg { "}}} notifyAll($value)" } } } class SimplePublisherImpl : SimpleEventImpl(), SimplePublisher { @Synchronized override fun publish(value: T) { dbg { "publish($value)" } notifyAll(value) } } class SimpleState(initialValue: T) : SimpleEventImpl() { var value: T = initialValue @Synchronized set(newValue) { dbg { "set($field -> $newValue)" } if (field != newValue) { field = newValue notifyAll(newValue) } } @Synchronized get @Synchronized override fun listen(listener: EventListener): AutoCloseable { dbg { "listener.notify(${listener.instanceName()}) -> $value" } listener.notify(value) return super.listen(listener) } } class Symbol(@JvmField val symbol: String) { override fun toString(): String = "<$symbol>" } val UNINITIALIZED: Any? = Symbol("UNINITIALIZED") fun combineSimpleEvents( a: SimpleEvent, b: SimpleEvent, transform: (T1, T2) -> R, ): SimpleEvent { return combineSimpleEventsInternal(arrayOf(a, b), { null }) { args: Array<*> -> @Suppress("UNCHECKED_CAST") transform(args[0] as T1, args[1] as T2) } } fun combineSimpleEvents( a: SimpleEvent, b: SimpleEvent, c: SimpleEvent, transform: (T1, T2, T3) -> R, ): SimpleEvent { return combineSimpleEventsInternal(arrayOf(a, b, c), { null }) { args: Array<*> -> @Suppress("UNCHECKED_CAST") transform(args[0] as T1, args[1] as T2, args[2] as T3) } } fun combineSimpleEventsInternal( inputEvents: Array>, arrayFactory: () -> Array?, transform: (Array) -> R, ): SimpleEvent { return object : SimpleEvent { override fun listen(listener: EventListener): AutoCloseable { val latestValues = arrayOfNulls(inputEvents.size) latestValues.fill(UNINITIALIZED) var remainingUninitializedValues = latestValues.size fun updateCombined() { val results = arrayFactory() @Suppress("UNCHECKED_CAST") val newValue = if (results == null) { transform(latestValues as Array) } else { (latestValues as Array).copyInto(results) transform(results as Array) } dbg { "listener.notify(${listener.instanceName()}) -> $newValue" } listener.notify(newValue) } val upstreamHandles = inputEvents.mapIndexed { i, event -> event.listen { value -> val prevValue = latestValues[i] if (prevValue == UNINITIALIZED) { remainingUninitializedValues-- } latestValues[i] = value if (remainingUninitializedValues == 0) { updateCombined() } } } return AutoCloseable { upstreamHandles.forEach { it.close() } } } } } internal inline fun SimpleEvent.transformSimpleObservable( crossinline block: EventListener.(T) -> Unit ): SimpleEvent { val upstream = this return object : SimpleEvent { override fun listen(listener: EventListener): AutoCloseable { return upstream.listen { value -> listener.block(value) } } } } internal inline fun SimpleEvent.map(crossinline transform: (T) -> R): SimpleEvent { return transformSimpleObservable { value -> notify(transform(value)) } } internal inline fun SimpleEvent.filter( crossinline predicate: (T) -> Boolean ): SimpleEvent { return transformSimpleObservable { value -> if (predicate(value)) { notify(value) } } } fun SimpleEvent.distinctUntilChanged(): SimpleEvent { val upstream = this return object : SimpleEvent { override fun listen(listener: EventListener): AutoCloseable { // Store previous value here so that the cache does not leak to other listeners var previousValue = UNINITIALIZED return upstream.listen { value -> if (previousValue === UNINITIALIZED || previousValue != value) { previousValue = value dbg { "listener.notify(${listener.instanceName()}) -> $value" } listener.notify(value) } } } } } fun SimpleEvent.sample(other: SimpleEvent, transform: (A, B) -> C): SimpleEvent { val upstream = this return object : SimpleEvent { override fun listen(listener: EventListener): AutoCloseable { val noVal = Any() val sampledRef = AtomicReference(noVal) val a = other.listen { sampledRef.set(it) } val b = upstream.listen { val sampled = sampledRef.get() if (sampled != noVal) { @Suppress("UNCHECKED_CAST") val transformedValue = transform(it, sampled as B) dbg { "listener.notify(${listener.instanceName()}) -> $transformedValue" } listener.notify(transformedValue) } } return AutoCloseable { a.close() b.close() } } } } interface SimpleSuspendableObserver { suspend fun awaitNextValue(): T fun cancel() } fun SimpleEventImpl.asSuspendableObserver(executor: Executor): SimpleSuspendableObserver { val activeContinuation = AtomicReference?>(null) listen { newValue -> executor.execute { val awaitingCont = activeContinuation.getAndSet(null) awaitingCont!!.resume(newValue) } } return object : SimpleSuspendableObserver { override suspend fun awaitNextValue(): T { return suspendCoroutine { c -> if (!activeContinuation.compareAndSet(null, c)) { fail("Only one awaiter permitted at a time.") } } } override fun cancel() { activeContinuation.getAndSet(null)?.resumeWithException(CancellationException()) } } } // Similar concept to SynchronousQueue; can only pass one value at a time, and can only pass // values if actively being listened to class SimpleSynchronousState() { var nextInput: Continuation? = null fun putValueOrThrow(newValue: T) { val c = nextInput if (c != null) { nextInput = null c.resume(newValue) } else { fail("No one is awaiting. Can't send new value if there are no listeners.") } } suspend fun awaitValue(): T { return suspendCoroutine { continuation -> if (nextInput != null) { fail("Already awaiting. Can't override next continuation") } nextInput = continuation } } } class SimpleWritableEventBuilder(val executor: Executor) : EventContext, WritableEventFactory>, EventContextProvider>, EventCombiner>, IntEventCombiner>, DistinctUntilChangedOperator>, FilterOperator>, MapOperator>, SampleOperator> { val writeContext = SimpleEventWriteContext() override fun createWritableEvent(value: T): SimpleObservableStateBoxOut { return SimpleEventBoxImpl(SimpleState(value)) } override fun combineEvents( a: SimpleEventBoxIn, b: SimpleEventBoxIn, transform: (T1, T2) -> R, ): SimpleObservableBoxOut { return SimpleEventBoxImpl(combineSimpleEvents(a.unbox(), b.unbox(), transform)) } override fun combineEvents( a: SimpleEventBoxIn, b: SimpleEventBoxIn, c: SimpleEventBoxIn, transform: (T1, T2, T3) -> R, ): SimpleObservableBoxOut { return SimpleEventBoxImpl(combineSimpleEvents(a.unbox(), b.unbox(), c.unbox(), transform)) } override fun combineIntEvents( events: Iterable>, transform: (Array) -> R, ): SimpleObservableBoxOut { val inputEvents = events.map { it.unbox() }.toList().toTypedArray() return SimpleEventBoxImpl( combineSimpleEventsInternal( inputEvents, { arrayOfNulls(inputEvents.size) }, transform, ) ) } override fun read(block: ReadContext>.() -> Unit): AutoCloseable { val readContext = SimpleEventObservationContext(executor) readContext.block() return readContext } override fun write(block: WriteContext>.() -> Unit) { writeContext.block() } override fun SimpleEventBoxIn.distinctUntilChanged(): SimpleObservableBoxOut { return SimpleEventBoxImpl(unbox().distinctUntilChanged()) } override fun SimpleEventBoxIn.filter( predicate: (T) -> Boolean ): SimpleObservableBoxOut { return SimpleEventBoxImpl(unbox().filter(predicate)) } override fun SimpleEventBoxIn.map(transform: (A) -> B): SimpleObservableBoxOut { this as SimpleEventBoxImpl return SimpleEventBoxImpl(unbox().map(transform)) } override fun SimpleEventBoxIn.sample( other: SimpleEventBoxIn, transform: (A, B) -> C, ): SimpleObservableBoxOut { return SimpleEventBoxImpl(unbox().sample(other.unbox(), transform)) } } class SimpleEventObservationContext(val executor: Executor) : ReadContext>, AutoCloseable { private val cancellationList = mutableListOf() private var closed = false override fun SimpleEventBoxIn.observe(block: (T) -> Unit) { val observable = unbox() // Start listening on bg thread, and always resume with latest value on the bg thread executor.execute { synchronized(cancellationList) { cancellationList += observable.listen { value -> executor.execute { block(value) } } } } } override fun close() { synchronized(cancellationList) { if (closed) { fail("Read context was already closed") } cancellationList.forEach { it.close() } closed = true } } } class SimpleEventWriteContext() : WriteContext> { override fun SimpleEventBoxIn.update(value: T) { (unbox() as SimpleState).value = value } override fun SimpleEventBoxIn.current(): T { return (unbox() as SimpleState).value } } abstract class BaseSimpleEventBenchmark(threadParam: ThreadFactory) : BaseEventBenchmark( threadParam, { SimpleWritableEventBuilder(it) }, )