/* * 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.systemui.util.kotlin.sample import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.MutableStateFlow import kotlinx.coroutines.flow.combine import kotlinx.coroutines.flow.distinctUntilChanged import kotlinx.coroutines.flow.filter import kotlinx.coroutines.flow.map import kotlinx.coroutines.launch typealias FlowBoxIn = EventBox> typealias FlowBoxOut = EventBox> typealias MutableStateFlowBoxOut = EventBox> private class FlowBoxImpl>(flow: R) : AbstractEventBox(flow) private fun FlowBoxIn.unbox(): Flow = (this as FlowBoxImpl).event class FlowWritableEventBuilder(val scope: CoroutineScope) : EventContext, WritableEventFactory>, EventContextProvider>, EventCombiner>, IntEventCombiner>, DistinctUntilChangedOperator>, FilterOperator>, MapOperator>, SampleOperator> { val writeContext = FlowWriteContext() override fun createWritableEvent(value: T): MutableStateFlowBoxOut { return FlowBoxImpl(MutableStateFlow(value)) } override fun combineEvents( a: FlowBoxIn, b: FlowBoxIn, transform: (T1, T2) -> T3, ): EventBox> { return FlowBoxImpl(combine(a.unbox(), b.unbox(), transform)) } override fun combineEvents( a: FlowBoxIn, b: FlowBoxIn, c: FlowBoxIn, transform: (T1, T2, T3) -> T4, ): FlowBoxOut { return FlowBoxImpl(combine(a.unbox(), b.unbox(), c.unbox(), transform)) } override fun combineIntEvents( events: Iterable>, transform: (Array) -> R, ): FlowBoxOut { val flows = events.map { it.unbox() } return FlowBoxImpl(combine(flows, transform)) } override fun read(block: ReadContext>.() -> Unit): AutoCloseable { val job = scope.launch { val readContext = FlowReadContext(this) readContext.block() } return AutoCloseable { job.cancel() } } override fun write(block: WriteContext>.() -> Unit) { writeContext.block() } override fun FlowBoxIn.distinctUntilChanged(): FlowBoxOut { return FlowBoxImpl(unbox().distinctUntilChanged()) } override fun FlowBoxIn.filter(predicate: (T) -> Boolean): FlowBoxOut { return FlowBoxImpl(unbox().filter(predicate)) } override fun FlowBoxIn.map(transform: (A) -> B): FlowBoxOut { return FlowBoxImpl(unbox().map(transform)) } override fun FlowBoxIn.sample( other: FlowBoxIn, transform: (A, B) -> C, ): FlowBoxOut { return FlowBoxImpl(unbox().sample(other.unbox(), transform)) } } class FlowReadContext(val scope: CoroutineScope) : ReadContext> { override fun FlowBoxIn.observe(block: (T) -> Unit) { scope.launch { unbox().collect(block) } } } class FlowWriteContext : WriteContext> { override fun FlowBoxIn.update(value: T) { (unbox() as MutableStateFlow).value = value } override fun FlowBoxIn.current(): T { return (unbox() as MutableStateFlow).value } } abstract class BaseFlowEventBenchmark(threadParam: ThreadFactory) : BaseEventBenchmark( threadParam, { FlowWritableEventBuilder(it) }, )