diff --git a/build-all.sh b/build-all.sh index a3c4772..489238b 100755 --- a/build-all.sh +++ b/build-all.sh @@ -16,8 +16,8 @@ cp clients/cli/build/libs/phoebench-cli.jar artifacts/ cp clients/compose/build/compose/binaries/main/app/com.jaytux.phoebench.clients/phoebench-linux.zip artifacts/ mkdir -p artifacts/wasm -cp clients/compose/build/kotlin-webpack/wasmJs/productionExecutable/* clients/compose/src/wasmJsMain/resources/* artifacts/wasm -(cd artifacts/wasm && zip phoebench-wasm.zip -r ./*) +cp -r clients/compose/build/kotlin-webpack/wasmJs/productionExecutable/* clients/compose/build/processedResources/wasmJs/main/* artifacts/wasm +(cd artifacts/wasm && rm config.json && zip phoebench-wasm.zip -r ./*) mv artifacts/wasm/phoebench-wasm.zip artifacts/ rm -fr artifacts/wasm diff --git a/build.gradle.kts b/build.gradle.kts index 693146e..7b66ab9 100644 --- a/build.gradle.kts +++ b/build.gradle.kts @@ -14,4 +14,4 @@ repositories { mavenCentral() } -version = PhoebenchVersion(1, 1, 0, "") \ No newline at end of file +version = PhoebenchVersion(1, 2, 0, "") \ No newline at end of file diff --git a/buildSrc/src/main/kotlin/BuildUtil.kt b/buildSrc/src/main/kotlin/BuildUtil.kt new file mode 100644 index 0000000..2c44bb5 --- /dev/null +++ b/buildSrc/src/main/kotlin/BuildUtil.kt @@ -0,0 +1,18 @@ +import org.gradle.api.Project + +fun Project.envValue(key: String): String? { + val fromEnv = providers.environmentVariable(key).orNull + if(fromEnv != null) return fromEnv + + val envFile = rootProject.file(".env") + if(!envFile.exists()) return null + + return envFile.useLines { lines -> + lines.map { it.trim() }.filter { it.isNotBlank() && !it.startsWith('#') && '=' in it } + .map { line -> + val (k, v) = line.split('=', limit = 2) + k.trim() to v.trim().removeSurrounding("\"").removeSurrounding("'") + } + .firstOrNull { it.first == key }?.second + } +} \ No newline at end of file diff --git a/clients/cli/build.gradle.kts b/clients/cli/build.gradle.kts index 83efef6..594ed62 100644 --- a/clients/cli/build.gradle.kts +++ b/clients/cli/build.gradle.kts @@ -5,10 +5,11 @@ plugins { application alias(libs.plugins.serialization) alias(libs.plugins.shadow) + alias(libs.plugins.buildconfig) } group = "com.jaytux.phoebench" -version = PhoebenchVersion(1, 1, 1) +version = rootProject.version as PhoebenchVersion if((version as PhoebenchVersion) < (rootProject.version as PhoebenchVersion)) throw GradleException("CLI Client version must be at least as high as protocol/common version") @@ -28,12 +29,15 @@ val generateVersion = tasks.register("serverVersion") { } } +val isDebug = envValue("PHOEBENCH_BUILD_RELEASE") == null + dependencies { implementation(kotlin("stdlib")) implementation(libs.clikt) implementation(libs.ktor.client.core) implementation(libs.ktor.client.auth) implementation(libs.ktor.client.content.negotiation) + implementation(libs.ktor.client.websocket) implementation(libs.kotlinx.datetime) implementation(libs.kotlinx.serialization) implementation(project(":common")) @@ -42,6 +46,7 @@ dependencies { implementation(libs.slf4j.simple) implementation(libs.java.keystore) implementation(libs.ktor.serialization.kotlinx.json) + implementation(libs.process) } application { @@ -75,5 +80,24 @@ kotlin { } tasks.withType { - archiveFileName = "phoebench-cli.jar" + archiveFileName = if(isDebug) "phoebench-cli-debug.jar" else "phoebench-cli.jar" +} + +buildConfig { + generateAtSync = false + useKotlinOutput { internalVisibility = false } + + println("Source sets: ${kotlin.sourceSets.toList().map { it.name }}") + + className("PersistenceConstants") + packageName("com.jaytux.phoebench.clients.cli") + + if(isDebug) { + println("Using debug CLI service") + buildConfigField("SERVICE", "com.jaytux.phoebench.cli.debug") + } + else { + println("Using release CLI service") + buildConfigField("SERVICE", "com.jaytux.phoebench.cli") + } } \ No newline at end of file diff --git a/clients/cli/src/main/kotlin/com/jaytux/phoebench/clients/cli/CLI.kt b/clients/cli/src/main/kotlin/com/jaytux/phoebench/clients/cli/CLI.kt index d03f3e6..2f9a637 100644 --- a/clients/cli/src/main/kotlin/com/jaytux/phoebench/clients/cli/CLI.kt +++ b/clients/cli/src/main/kotlin/com/jaytux/phoebench/clients/cli/CLI.kt @@ -1,6 +1,8 @@ package com.jaytux.phoebench.clients.cli import com.github.ajalt.clikt.core.* +import com.github.ajalt.clikt.parameters.arguments.argument +import com.github.ajalt.clikt.parameters.arguments.multiple import com.github.ajalt.clikt.parameters.groups.mutuallyExclusiveOptions import com.github.ajalt.clikt.parameters.groups.single import com.github.ajalt.clikt.parameters.options.* @@ -227,6 +229,24 @@ object CLI { } } + @Suppress("unused") + class Monitor : CliktCommand(name = "monitor") { + val finder by mutuallyExclusiveOptions( + option("--id", help = "Find a project by UUID.").convert { Project.ID(Uuid.parse(it)) }, + option("--name", help = "Find a project by name (formatted [user]/[project])").convert { + val split = it.split('/') + if(split.size != 2) throw IllegalArgumentException("Invalid format (expected [user]/[project])") + Project.ProjectName(split[0], split[1]) + } + ).single() + val command by argument("command", help = "The command to be run") + val commandArgs by argument("arguments", help = "Arguments to pass to the command").multiple() + + override fun run() { + MonitorHandler.monitor(finder, command, commandArgs) + } + } + @Suppress("unused") class Version : CliktCommand(name = "version") { override fun run() { diff --git a/clients/cli/src/main/kotlin/com/jaytux/phoebench/clients/cli/Client.kt b/clients/cli/src/main/kotlin/com/jaytux/phoebench/clients/cli/Client.kt index 5aa34c3..9698d34 100644 --- a/clients/cli/src/main/kotlin/com/jaytux/phoebench/clients/cli/Client.kt +++ b/clients/cli/src/main/kotlin/com/jaytux/phoebench/clients/cli/Client.kt @@ -1,6 +1,7 @@ package com.jaytux.phoebench.clients.cli import com.jaytux.phoebench.common.ApiRoute +import com.jaytux.phoebench.common.CSERoute import com.jaytux.phoebench.common.Either import com.jaytux.phoebench.common.ErrorResponse import com.jaytux.phoebench.common.IClient @@ -9,15 +10,31 @@ import com.jaytux.phoebench.common.Routes import com.jaytux.phoebench.common.TokenResponse import com.jaytux.phoebench.common.error import com.jaytux.phoebench.common.foldSuspend +import com.jaytux.phoebench.common.value import io.ktor.client.HttpClient +import io.ktor.client.call.body import io.ktor.client.engine.okhttp.OkHttp +import io.ktor.client.plugins.ResponseException import io.ktor.client.plugins.auth.Auth import io.ktor.client.plugins.auth.providers.BearerTokens import io.ktor.client.plugins.auth.providers.bearer import io.ktor.client.plugins.contentnegotiation.ContentNegotiation +import io.ktor.client.plugins.websocket.WebSocketException +import io.ktor.client.plugins.websocket.WebSockets +import io.ktor.client.plugins.websocket.sendSerialized +import io.ktor.client.plugins.websocket.webSocket +import io.ktor.client.request.url +import io.ktor.http.HttpMethod +import io.ktor.serialization.kotlinx.KotlinxWebsocketSerializationConverter import io.ktor.serialization.kotlinx.json.json +import io.ktor.util.reflect.TypeInfo +import io.ktor.util.reflect.typeInfo import io.ktor.utils.io.CancellationException +import io.ktor.websocket.CloseReason import kotlinx.coroutines.asExecutor +import kotlinx.serialization.KSerializer +import kotlinx.serialization.json.Json +import java.net.URL import kotlin.uuid.Uuid object Client { @@ -68,6 +85,10 @@ object Client { } } } + + install(WebSockets) { + contentConverter = KotlinxWebsocketSerializationConverter(Json) + } } init { @@ -91,6 +112,19 @@ object Client { suspend fun callRoute(route: ApiRoute, body: TReq): Either = callRoute(_client, route, body) + suspend fun callCSE( + route: CSERoute, params: TParams, + body: suspend (sender: suspend (TEvent) -> Unit) -> Unit + ): Either { + try { + val client = IClient.Default(_client, _server ?: throw IllegalStateException("No server URL set.")) + return route.call(client, params, body) + } + catch(e: CancellationException) { + return ErrorResponse("Event stream disconnected.").error() + } + } + fun onLogin(tokens: TokenResponse) { _refreshToken = tokens.refresh _accessToken = tokens.access diff --git a/clients/cli/src/main/kotlin/com/jaytux/phoebench/clients/cli/MonitorHandler.kt b/clients/cli/src/main/kotlin/com/jaytux/phoebench/clients/cli/MonitorHandler.kt new file mode 100644 index 0000000..c1e49f2 --- /dev/null +++ b/clients/cli/src/main/kotlin/com/jaytux/phoebench/clients/cli/MonitorHandler.kt @@ -0,0 +1,83 @@ +package com.jaytux.phoebench.clients.cli + +import com.github.pgreze.process.Redirect +import com.github.pgreze.process.process +import com.jaytux.phoebench.clients.cli.ProjectHandlers.toId +import com.jaytux.phoebench.common.ANSI +import com.jaytux.phoebench.common.ClientMonitorEvent +import com.jaytux.phoebench.common.Routes +import com.jaytux.phoebench.common.Stream +import com.jaytux.phoebench.common.bind +import com.jaytux.phoebench.common.fold +import io.ktor.http.HttpMethod +import io.ktor.utils.io.CancellationException +import kotlin.random.Random +import kotlin.time.Clock +import kotlin.time.Instant + +object MonitorHandler { + private fun filterAnsi(line: String): Pair> { + val builder = StringBuilder() + val codes = mutableSetOf() + var remaining = line + while(remaining.isNotEmpty()) { + val next = remaining.indexOf("\u001B[") + if(next < 0) { + builder.append(remaining) + remaining = "" + } + else { + builder.append(remaining.substring(0, next)) + val end = remaining.indexOf('m', startIndex = next + 2) + if(end == -1) return (line to listOf()) // invalid... + codes += remaining.substring(next + 2, end).split(';').mapNotNull { + it.toIntOrNull()?.let { i -> ANSI.ansiMapping[i] } + } + remaining = remaining.substring(end + 1) + } + } + + return builder.toString() to codes.toList() + } + + fun monitor(find: CLI.Commands.Project.IProjectIdentification?, command: String, args: List) { + val project = ProjectHandlers.ensureProjectIdentification(find) + val combinedCommand = arrayOf(command, *args.toTypedArray()) + tryAuthenticated { + project.toId().bind { id -> + Client.callCSE(Routes.CSE.monitor, id) { sender -> + try { + val res = process( + *combinedCommand, + stdin = null, + stdout = Redirect.Consume { flow -> + flow.collect { line -> + val (updLine, options) = filterAnsi(line) + sender(ClientMonitorEvent.Message(msg = updLine, stream = Stream.STDOUT, options = options)) + println(line) + } + }, + stderr = Redirect.Consume { flow -> + flow.collect { line -> + val (updLine, options) = filterAnsi(line) + sender(ClientMonitorEvent.Message(msg = updLine, stream = Stream.STDERR, options = options)) + System.err.println(line) + } + } + ) + + sender(ClientMonitorEvent.ApplicationFinished(res.resultCode)) + } + catch(e: CancellationException) { + throw e + } + catch(e: Exception) { + System.err.println("Launch failed: ${e.message} (${e::class.simpleName})") + } + } + } + }.fold({ + System.err.println("Failed to run: ${it.msg}") + }) {} + } +} \ No newline at end of file diff --git a/clients/cli/src/main/kotlin/com/jaytux/phoebench/clients/cli/PersistentStorage.kt b/clients/cli/src/main/kotlin/com/jaytux/phoebench/clients/cli/PersistentStorage.kt index df149e3..e1bb9b6 100644 --- a/clients/cli/src/main/kotlin/com/jaytux/phoebench/clients/cli/PersistentStorage.kt +++ b/clients/cli/src/main/kotlin/com/jaytux/phoebench/clients/cli/PersistentStorage.kt @@ -10,7 +10,7 @@ import kotlin.uuid.Uuid object PersistentStorage { private val json = Json - const val SERVICE = "com.jaytux.phoebench.cli" + const val SERVICE = PersistenceConstants.SERVICE private var _backend: IBackend = KeyringBackend interface IBackend { diff --git a/clients/compose/build.gradle.kts b/clients/compose/build.gradle.kts index b62aaf9..3f9c9ea 100644 --- a/clients/compose/build.gradle.kts +++ b/clients/compose/build.gradle.kts @@ -13,7 +13,7 @@ plugins { alias(libs.plugins.buildconfig) } -version = PhoebenchVersion(1, 1, 2) +version = rootProject.version as PhoebenchVersion if((version as PhoebenchVersion) < (rootProject.version as PhoebenchVersion)) throw GradleException("UI Clients version must be at least as high as protocol/common version") @@ -142,23 +142,6 @@ compose.desktop { } } -fun envValue(key: String): String? { - val fromEnv = providers.environmentVariable(key).orNull - if(fromEnv != null) return fromEnv - - val envFile = rootProject.file(".env") - if(!envFile.exists()) return null - - return envFile.useLines { lines -> - lines.map { it.trim() }.filter { it.isNotBlank() && !it.startsWith('#') && '=' in it } - .map { line -> - val (k, v) = line.split('=', limit = 2) - k.trim() to v.trim().removeSurrounding("\"").removeSurrounding("'") - } - .firstOrNull { it.first == key }?.second - } -} - buildConfig { generateAtSync = false useKotlinOutput { internalVisibility = false } diff --git a/clients/compose/src/commonMain/composeResources/font/JetBrainsMono-Bold.ttf b/clients/compose/src/commonMain/composeResources/font/JetBrainsMono-Bold.ttf new file mode 100644 index 0000000..8c93043 Binary files /dev/null and b/clients/compose/src/commonMain/composeResources/font/JetBrainsMono-Bold.ttf differ diff --git a/clients/compose/src/commonMain/composeResources/font/JetBrainsMono-BoldItalic.ttf b/clients/compose/src/commonMain/composeResources/font/JetBrainsMono-BoldItalic.ttf new file mode 100644 index 0000000..1ddf216 Binary files /dev/null and b/clients/compose/src/commonMain/composeResources/font/JetBrainsMono-BoldItalic.ttf differ diff --git a/clients/compose/src/commonMain/composeResources/font/JetBrainsMono-Italic.ttf b/clients/compose/src/commonMain/composeResources/font/JetBrainsMono-Italic.ttf new file mode 100644 index 0000000..ccc9d6a Binary files /dev/null and b/clients/compose/src/commonMain/composeResources/font/JetBrainsMono-Italic.ttf differ diff --git a/clients/compose/src/commonMain/composeResources/font/JetBrainsMono-Regular.ttf b/clients/compose/src/commonMain/composeResources/font/JetBrainsMono-Regular.ttf new file mode 100644 index 0000000..dff66cc Binary files /dev/null and b/clients/compose/src/commonMain/composeResources/font/JetBrainsMono-Regular.ttf differ diff --git a/clients/compose/src/commonMain/kotlin/com/jaytux/phoebench/clients/PlatformAPI.kt b/clients/compose/src/commonMain/kotlin/com/jaytux/phoebench/clients/PlatformAPI.kt index 24af67d..24597b5 100644 --- a/clients/compose/src/commonMain/kotlin/com/jaytux/phoebench/clients/PlatformAPI.kt +++ b/clients/compose/src/commonMain/kotlin/com/jaytux/phoebench/clients/PlatformAPI.kt @@ -4,8 +4,7 @@ import androidx.compose.runtime.Composable import androidx.compose.ui.draganddrop.DragAndDropEvent import androidx.compose.ui.draganddrop.DragAndDropTransferData import androidx.compose.ui.platform.ClipEntry -import io.ktor.client.HttpClient -import io.ktor.client.HttpClientConfig +import io.ktor.client.* import kotlinx.serialization.KSerializer import kotlinx.serialization.serializer import kotlin.uuid.Uuid diff --git a/clients/compose/src/commonMain/kotlin/com/jaytux/phoebench/clients/Util.kt b/clients/compose/src/commonMain/kotlin/com/jaytux/phoebench/clients/Util.kt index 2a1a5e1..df6e943 100644 --- a/clients/compose/src/commonMain/kotlin/com/jaytux/phoebench/clients/Util.kt +++ b/clients/compose/src/commonMain/kotlin/com/jaytux/phoebench/clients/Util.kt @@ -1,9 +1,15 @@ package com.jaytux.phoebench.clients +import androidx.compose.material3.Typography +import androidx.compose.runtime.Composable import androidx.compose.runtime.MutableState import androidx.compose.runtime.State import androidx.compose.ui.graphics.Color import androidx.compose.ui.graphics.toArgb +import androidx.compose.ui.text.TextStyle +import androidx.compose.ui.text.font.FontFamily +import androidx.compose.ui.text.font.FontStyle +import androidx.compose.ui.text.font.FontWeight import androidx.lifecycle.ViewModel import androidx.lifecycle.viewModelScope import com.jaytux.phoebench.common.Either @@ -18,9 +24,10 @@ import kotlinx.datetime.format.MonthNames import kotlinx.datetime.format.Padding import kotlinx.datetime.format.char import kotlinx.datetime.toLocalDateTime +import org.jetbrains.compose.resources.Font +import phoebench.clients.compose.generated.resources.* import kotlin.math.absoluteValue import kotlin.math.pow -import kotlin.math.roundToInt import kotlin.random.Random import kotlin.random.nextInt import kotlin.time.Clock @@ -144,4 +151,16 @@ fun List.geomean(): Float { var prod = 1.0f for(value in this) prod *= value return prod.pow(1.0f / size.toFloat()) +} + +@Composable +fun TextStyle.makeMonospaced(): TextStyle { + val family = FontFamily( + Font(Res.font.JetBrainsMono_Regular, weight = FontWeight.Normal, style = FontStyle.Normal), + Font(Res.font.JetBrainsMono_Bold, weight = FontWeight.Bold, style = FontStyle.Normal), + Font(Res.font.JetBrainsMono_Italic, weight = FontWeight.Normal, style = FontStyle.Italic), + Font(Res.font.JetBrainsMono_BoldItalic, weight = FontWeight.Bold, style = FontStyle.Italic) + ) + + return copy(fontFamily = family) } \ No newline at end of file diff --git a/clients/compose/src/commonMain/kotlin/com/jaytux/phoebench/clients/data/IProjectRepo.kt b/clients/compose/src/commonMain/kotlin/com/jaytux/phoebench/clients/data/IProjectRepo.kt index 572e857..80fc69b 100644 --- a/clients/compose/src/commonMain/kotlin/com/jaytux/phoebench/clients/data/IProjectRepo.kt +++ b/clients/compose/src/commonMain/kotlin/com/jaytux/phoebench/clients/data/IProjectRepo.kt @@ -38,6 +38,8 @@ interface IProjectRepo { unit: TimeUnit, input: String, hardware: String): Either suspend fun deleteEntry(id: Uuid): Either + suspend fun deleteMonitorLogs(): Either + companion object { class Default(private val _client: Client, private val _projectId: Uuid) : IProjectRepo { override suspend fun get(): Either = @@ -74,6 +76,9 @@ interface IProjectRepo { override suspend fun deleteEntry(id: Uuid): Either = _client.callRoute(Routes.Entry.delete, id).ignoreValue() + + override suspend fun deleteMonitorLogs(): Either = + _client.callRoute(Routes.Project.rmMonitor, _projectId).ignoreValue() } fun default(client: Client, projectId: Uuid) = Default(client, projectId) diff --git a/clients/compose/src/commonMain/kotlin/com/jaytux/phoebench/clients/data/ISSERepo.kt b/clients/compose/src/commonMain/kotlin/com/jaytux/phoebench/clients/data/ISSERepo.kt index bde2e87..81f995a 100644 --- a/clients/compose/src/commonMain/kotlin/com/jaytux/phoebench/clients/data/ISSERepo.kt +++ b/clients/compose/src/commonMain/kotlin/com/jaytux/phoebench/clients/data/ISSERepo.kt @@ -7,12 +7,14 @@ import com.jaytux.phoebench.common.ErrorResponse import com.jaytux.phoebench.common.HomeEvent import com.jaytux.phoebench.common.ProjectEvent import com.jaytux.phoebench.common.Routes +import com.jaytux.phoebench.common.ServerMonitorEvent import kotlin.uuid.Uuid interface ISSERepo { suspend fun connectHome(onEvent: suspend (Either) -> Unit): Either suspend fun connectAdmin(onEvent: suspend (Either) -> Unit): Either suspend fun connectProject(id: Uuid, onEvent: suspend (Either) -> Unit): Either + suspend fun connectMonitor(id: Uuid, onEvent: suspend (Either) -> Unit): Either companion object { class Default(private val _client: Client) : ISSERepo { @@ -24,6 +26,9 @@ interface ISSERepo { override suspend fun connectProject(id: Uuid, onEvent: suspend (Either) -> Unit): Either = _client.callSSE(Routes.SSE.projectSpecific, id, onEvent) + + override suspend fun connectMonitor(id: Uuid, onEvent: suspend (Either) -> Unit): Either = + _client.callSSE(Routes.SSE.monitor, id, onEvent) } fun default(client: Client) = Default(client) diff --git a/clients/compose/src/commonMain/kotlin/com/jaytux/phoebench/clients/data/ProjectVM.kt b/clients/compose/src/commonMain/kotlin/com/jaytux/phoebench/clients/data/ProjectVM.kt index 9165dcb..10d4551 100644 --- a/clients/compose/src/commonMain/kotlin/com/jaytux/phoebench/clients/data/ProjectVM.kt +++ b/clients/compose/src/commonMain/kotlin/com/jaytux/phoebench/clients/data/ProjectVM.kt @@ -1,5 +1,6 @@ package com.jaytux.phoebench.clients.data +import androidx.compose.runtime.mutableStateListOf import androidx.compose.runtime.mutableStateOf import androidx.compose.ui.graphics.Color import androidx.lifecycle.ViewModel @@ -18,7 +19,9 @@ import com.jaytux.phoebench.common.BenchmarkResponse import com.jaytux.phoebench.common.EntryResponse import com.jaytux.phoebench.common.VersionResponse import com.jaytux.phoebench.common.ProjectEvent +import com.jaytux.phoebench.common.ServerMonitorEvent import com.jaytux.phoebench.common.TimeUnit +import com.jaytux.phoebench.common.fold import kotlinx.coroutines.Job import kotlin.time.Clock import kotlin.time.Instant @@ -87,6 +90,7 @@ class ProjectVM( private val _currentBenchmark = mutableStateOf(0) private val _inputs = mutableStateOf(setOf()) private val _hardware = mutableStateOf(setOf()) + private val _monitorMessages = mutableStateOf?>(null) val name = _name.immutable() val owner = _owner.immutable() @@ -97,8 +101,10 @@ class ProjectVM( val currentBenchmark = _currentBenchmark.immutable() val inputs = _inputs.immutable() val hardware = _hardware.immutable() + val monitorMessages = _monitorMessages.immutable() private var _job: Job? = null + private var _monitorJob: Job? = null init { _job = withScope { @@ -221,7 +227,7 @@ class ProjectVM( val (updated, idx) = _benchmarks.value.insortIdx( Benchmark.fromResponse(bench, _labels.value, { i -> _inputs.value += i }, { hw -> _hardware.value += hw }) ) { it.name } - _currentBenchmark.value?.let { curr -> + _currentBenchmark.value.let { curr -> if(idx != null) { if(idx < curr) _currentBenchmark.value = curr + 1 } @@ -279,4 +285,37 @@ class ProjectVM( fun selectBenchmark(id: Uuid) { _currentBenchmark.value = maxOf(_benchmarks.value.indexOfFirst { it.id == id }, 0) } + + fun openMonitor() { + if(_monitorJob != null) return + _monitorMessages.value = listOf() + _monitorJob = withScope { + _sseRepo.connectMonitor(_id) { ev -> + ev.snackOr { + when(it) { + is ServerMonitorEvent.Backlog -> { + val res = ArrayList(it.messages.size + 2) + it.start?.let { start -> res += start } + res.addAll(it.messages) + it.end?.let { end -> res += end } + _monitorMessages.value = res + } + ServerMonitorEvent.Cleared -> _monitorMessages.value = listOf() + is ServerMonitorEvent.ApplicationEnd, is ServerMonitorEvent.ApplicationStart, is ServerMonitorEvent.Message -> + _monitorMessages.value = (_monitorMessages.value ?: listOf()) + it + } + } + } + } + } + + fun closeMonitor() { + _monitorJob?.cancel() + _monitorJob = null + _monitorMessages.value = null + } + + fun clearMonitor() = withScope { + _repo.deleteMonitorLogs().snackOnError() + } } \ No newline at end of file diff --git a/clients/compose/src/commonMain/kotlin/com/jaytux/phoebench/clients/ui/ProjectView.kt b/clients/compose/src/commonMain/kotlin/com/jaytux/phoebench/clients/ui/ProjectView.kt index c87c11e..6233ce1 100644 --- a/clients/compose/src/commonMain/kotlin/com/jaytux/phoebench/clients/ui/ProjectView.kt +++ b/clients/compose/src/commonMain/kotlin/com/jaytux/phoebench/clients/ui/ProjectView.kt @@ -18,12 +18,14 @@ import androidx.compose.ui.graphics.Color import androidx.compose.ui.graphics.SolidColor import androidx.compose.ui.layout.onGloballyPositioned import androidx.compose.ui.platform.LocalDensity +import androidx.compose.ui.platform.LocalTextToolbar import androidx.compose.ui.text.font.FontStyle import androidx.compose.ui.text.font.FontWeight import androidx.compose.ui.text.style.TextDecoration import androidx.compose.ui.text.style.TextOverflow import androidx.compose.ui.unit.dp import androidx.compose.ui.window.Dialog +import androidx.compose.ui.window.DialogProperties import androidx.lifecycle.viewmodel.compose.viewModel import com.composables.icons.lucide.* import com.jaytux.phoebench.clients.* @@ -57,6 +59,7 @@ fun ProjectView(id: Uuid, forceBack: () -> Unit) { val versions by vm.versions val benchmarks by vm.benchmarks val currentBenchmark by vm.currentBenchmark + val monitor by vm.monitorMessages var editing by remember { mutableStateOf(false) } var deleting by remember { mutableStateOf(false) } @@ -64,16 +67,20 @@ fun ProjectView(id: Uuid, forceBack: () -> Unit) { Column(Modifier.padding(all = 15.dp)) { Row(Modifier.height(IntrinsicSize.Min), verticalAlignment = Alignment.CenterVertically) { - Text("Project ${name ?: "Unnamed Project"}", style = MaterialTheme.typography.headlineMedium) - if(editable) { - Spacer(Modifier.width(25.dp)) - IconButton({ editing = true }) { - Icon(Lucide.Pencil, "Edit project details") + Text("Project ${name ?: "Unnamed Project"}", style = MaterialTheme.typography.headlineMedium) + if (editable) { + Spacer(Modifier.width(25.dp)) + IconButton({ editing = true }) { + Icon(Lucide.Pencil, "Edit project details") + } + IconButton({ deleting = true }) { + Icon(Lucide.Trash, "Delete project") + } + IconButton({ vm.openMonitor() }) { + Icon(Lucide.SquareTerminal, "Open monitor") + } } - IconButton({ deleting = true }) { - Icon(Lucide.Trash, "Delete project") - } - } + } owner?.let { Text("${if(public) "Public" else "Private"} project by user $it") } Spacer(Modifier.height(15.dp)) @@ -112,6 +119,82 @@ fun ProjectView(id: Uuid, forceBack: () -> Unit) { vm.mkBenchmark(it) addingBenchmark = false } + + monitor?.let { MonitorDialog(name ?: "Unnamed Project", it, vm::closeMonitor, vm::clearMonitor) } +} + +@Composable +fun MonitorMessage(msg: ServerMonitorEvent.ITextEvent) = when(msg) { + is ServerMonitorEvent.ApplicationEnd -> Text("[${msg.time.fmt()}] Application finished with exit code ${msg.exitCode}.", fontStyle = FontStyle.Italic) + is ServerMonitorEvent.ApplicationStart -> Text("[${msg.time.fmt()}] Application started.", fontStyle = FontStyle.Italic) + is ServerMonitorEvent.Message -> { + val ansi = ANSI.unpack(msg.options).toSet() + + val weight = if(ANSI.BOLD in ansi) FontWeight.Bold else null + val decoration = if(ANSI.UNDERLINE in ansi) TextDecoration.Underline else null + val color = when { + ANSI.FG_BLACK in ansi -> Color.Black + ANSI.FG_RED in ansi -> Color.Red + ANSI.FG_GREEN in ansi -> Color.Green + ANSI.FG_YELLOW in ansi -> Color.Yellow + ANSI.FG_BLUE in ansi -> Color.Blue + ANSI.FG_MAGENTA in ansi -> Color.Magenta + ANSI.FG_CYAN in ansi -> Color.Cyan + else -> LocalContentColor.current + } + + val background = when { + ANSI.BG_BLACK in ansi -> Color.Black + ANSI.BG_RED in ansi -> Color.Red + ANSI.BG_GREEN in ansi -> Color.Green + ANSI.BG_YELLOW in ansi -> Color.Yellow + ANSI.BG_BLUE in ansi -> Color.Blue + ANSI.BG_MAGENTA in ansi -> Color.Magenta + ANSI.BG_CYAN in ansi -> Color.Cyan + else -> null + }?.let { Modifier.background(it) } ?: Modifier + + Box(background.fillMaxWidth()) { + Text( + "[${msg.timeStamp.fmt()}] [${if (msg.stream == Stream.STDOUT) 'O' else 'E'}] ${msg.message}", + fontWeight = weight, textDecoration = decoration, color = color + ) + } + } +} + +@Composable +fun MonitorDialog(name: String, logs: List, onClose: () -> Unit, onClear: () -> Unit) { + Dialog(onDismissRequest = onClose, properties = DialogProperties(usePlatformDefaultWidth = false)) { + Surface(Modifier.padding(15.dp).fillMaxHeight(0.8f).widthIn(min = 300.dp, max = 2000.dp), shape = MaterialTheme.shapes.medium) { + Column(Modifier.padding(8.dp)) { + Text("Monitor for $name", Modifier.align(Alignment.CenterHorizontally), style = MaterialTheme.typography.headlineMedium) + Spacer(Modifier.height(10.dp)) + if(logs.isEmpty()) { + Box(Modifier.weight(1f)) { + Box(Modifier.fillMaxHeight(0.33f).fillMaxWidth()) { + Text("No monitor logs.", Modifier.align(Alignment.Center), fontStyle = FontStyle.Italic) + } + } + } + else { + CompositionLocalProvider(LocalTextStyle provides LocalTextStyle.current.makeMonospaced()) { + LazyColumn(Modifier.weight(1f).padding(10.dp).background(MaterialTheme.colorScheme.surfaceDim)) { + items(logs) { + MonitorMessage(it) + } + } + } + } + Spacer(Modifier.height(10.dp)) + Row { + Button(onClear, Modifier.weight(0.5f)) { Text("Clear Monitor") } + Spacer(Modifier.width(10.dp)) + Button(onClose, Modifier.weight(0.5f)) { Text("Close Monitor") } + } + } + } + } } @Composable diff --git a/clients/compose/src/desktopMain/kotlin/com/jaytux/phoebench/clients/PlatformAPI.desktop.kt b/clients/compose/src/desktopMain/kotlin/com/jaytux/phoebench/clients/PlatformAPI.desktop.kt index 86af3da..7486d2f 100644 --- a/clients/compose/src/desktopMain/kotlin/com/jaytux/phoebench/clients/PlatformAPI.desktop.kt +++ b/clients/compose/src/desktopMain/kotlin/com/jaytux/phoebench/clients/PlatformAPI.desktop.kt @@ -2,11 +2,7 @@ package com.jaytux.phoebench.clients import androidx.compose.runtime.Composable import androidx.compose.ui.ExperimentalComposeUiApi -import androidx.compose.ui.draganddrop.DragAndDropEvent -import androidx.compose.ui.draganddrop.DragAndDropTransferAction -import androidx.compose.ui.draganddrop.DragAndDropTransferData -import androidx.compose.ui.draganddrop.DragAndDropTransferable -import androidx.compose.ui.draganddrop.awtTransferable +import androidx.compose.ui.draganddrop.* import androidx.compose.ui.platform.ClipEntry import com.github.javakeyring.Keyring import com.jaytux.phoebench.clients.ui.DefaultServerSelect diff --git a/clients/compose/src/wasmJsMain/kotlin/com/jaytux/phoebench/clients/PlatformAPI.wasmJs.kt b/clients/compose/src/wasmJsMain/kotlin/com/jaytux/phoebench/clients/PlatformAPI.wasmJs.kt index e8239ff..4853cd0 100644 --- a/clients/compose/src/wasmJsMain/kotlin/com/jaytux/phoebench/clients/PlatformAPI.wasmJs.kt +++ b/clients/compose/src/wasmJsMain/kotlin/com/jaytux/phoebench/clients/PlatformAPI.wasmJs.kt @@ -13,8 +13,6 @@ import kotlinx.browser.window import kotlinx.coroutines.await import kotlinx.serialization.KSerializer import kotlinx.serialization.Serializable -import kotlinx.serialization.decodeFromString -import kotlinx.serialization.encodeToString import kotlinx.serialization.json.Json import kotlinx.serialization.serializer import org.w3c.dom.DataTransfer diff --git a/clients/compose/src/wasmJsMain/resources/index.html b/clients/compose/src/wasmJsMain/resources/index.html index 1567ed2..aa474ef 100644 --- a/clients/compose/src/wasmJsMain/resources/index.html +++ b/clients/compose/src/wasmJsMain/resources/index.html @@ -1,10 +1,14 @@ - - + + PhoeBench - + + + + + diff --git a/common/build.gradle.kts b/common/build.gradle.kts index 4051413..6601556 100644 --- a/common/build.gradle.kts +++ b/common/build.gradle.kts @@ -18,7 +18,7 @@ val generatePartials = tasks.register("generatePartials") { description = "Generate Partial classes (requests with all-nullable fields)" val scriptFile = project.file("partialize.main.kts") val targets = fileTree(requestsDirectory) { - include("**/*.kt") + include("Requests.kt") } val lst = targets.map { it.absolutePath } diff --git a/common/src/commonMain/kotlin/com/jaytux/phoebench/common/CSERoute.kt b/common/src/commonMain/kotlin/com/jaytux/phoebench/common/CSERoute.kt new file mode 100644 index 0000000..fd22149 --- /dev/null +++ b/common/src/commonMain/kotlin/com/jaytux/phoebench/common/CSERoute.kt @@ -0,0 +1,62 @@ +package com.jaytux.phoebench.common + +import io.ktor.client.plugins.websocket.* +import io.ktor.http.* +import io.ktor.util.reflect.* +import io.ktor.utils.io.* +import io.ktor.websocket.* +import kotlin.uuid.Uuid + +sealed class CSERoute(val path: String, val elevation: Elevation, private val _eventType: TypeInfo) { + open val pattern = path + + protected open fun buildUrl(params: TParams): String = path + + abstract fun extractParams(reqParams: Parameters): TParams? + + suspend fun call(client: IClient, params: TParams, body: suspend (sender: suspend (TEvent) -> Unit) -> Unit): Either { + val server = client.serverUrl.replace("https", "ws").replace("http", "ws") + var error: ErrorResponse? = null + try { + client.client.webSocket("$server${buildUrl(params)}", {}) { + try { + body { sendSerialized(it, _eventType) } + } + catch(e: CancellationException) { + closeReason.await()?.let { + if(it.code != CloseReason.Codes.NORMAL.code) error = ErrorResponse(it.message) + } + throw e + } + } + return Unit.value() + } + catch(e: CancellationException) { + return (error ?: ErrorResponse(e.message ?: "Unknown websocket error")).error() + } + catch(e: WebSocketException) { + return ErrorResponse("Could not set up websocket stream: ${e.message}").error() + } + catch(e: Exception) { + return ErrorResponse("WebSocket connection failed: ${e.message}").error() + } + } + + class CSERoute1(path: String, elevation: Elevation, eventType: TypeInfo, val urlEncode: (T1) -> String, val urlDecode: (String?) -> T1?) + : CSERoute(path, elevation, eventType) + { + override val pattern: String = "$path/{param}" + override fun buildUrl(params: T1): String = "$path/${urlEncode(params)}" + override fun extractParams(reqParams: Parameters): T1? = urlDecode(reqParams["param"]) + } + + companion object { + inline fun single(path: String, elevation: Elevation, + noinline urlEncode: (T1) -> String = { it.toString() }, noinline urlDecode: (String?) -> T1? + ) = CSERoute1(path, elevation, typeInfo(), urlEncode, urlDecode) + + inline fun uuid(path: String, elevation: Elevation) = single(path, elevation) { + it?.let { p -> Uuid.parseOrNull(p) } + } + } +} \ No newline at end of file diff --git a/common/src/commonMain/kotlin/com/jaytux/phoebench/common/CloseReasons.kt b/common/src/commonMain/kotlin/com/jaytux/phoebench/common/CloseReasons.kt new file mode 100644 index 0000000..3c2e142 --- /dev/null +++ b/common/src/commonMain/kotlin/com/jaytux/phoebench/common/CloseReasons.kt @@ -0,0 +1,8 @@ +package com.jaytux.phoebench.common + +enum class CloseReasons(val code: Short) { + NOT_AUTHORIZED(4001), + INVALID_REQUEST(4002), + CONFLICT(4003), + NOT_FOUND(4004) +} \ No newline at end of file diff --git a/common/src/commonMain/kotlin/com/jaytux/phoebench/common/Events.kt b/common/src/commonMain/kotlin/com/jaytux/phoebench/common/Events.kt index cc06f5d..232d6bf 100644 --- a/common/src/commonMain/kotlin/com/jaytux/phoebench/common/Events.kt +++ b/common/src/commonMain/kotlin/com/jaytux/phoebench/common/Events.kt @@ -1,6 +1,8 @@ package com.jaytux.phoebench.common import kotlinx.serialization.Serializable +import kotlin.collections.fold +import kotlin.time.Instant import kotlin.uuid.Uuid @Serializable @@ -65,4 +67,83 @@ sealed class ProjectEvent { @Serializable data class EntryDeleted(val id: Uuid, val benchmarkId: Uuid) : ProjectEvent() +} + +@Serializable + enum class Stream { + STDOUT, STDERR +} + +@Serializable +enum class ANSI(val bitIdx: Int, val ansiCode: Int) { + BOLD(0, 1), UNDERLINE(1, 4), + + FG_BLACK(2, 30), FG_RED(3, 31), FG_GREEN(4, 32), + FG_YELLOW(5, 33), FG_BLUE(6, 34), FG_MAGENTA(7, 35), + FG_CYAN(8, 36), + + BG_BLACK(9, 40), BG_RED(10, 41), BG_GREEN(11, 42), + BG_YELLOW(12, 43), BG_BLUE(13, 44), BG_MAGENTA(14, 45), + BG_CYAN(15, 46); + + companion object { + private val _mapping: Map + val ansiMapping: Map + init { + if(entries.map { it.bitIdx }.toSet().size != entries.size) + throw IllegalStateException("ANSI bit indices contain duplicates") + if(entries.map { it.ansiCode }.toSet().size != entries.size) + throw IllegalStateException("ANSI escape codes contain duplicates") + + _mapping = entries.associateBy { ansi -> ansi.bitIdx } + ansiMapping = entries.associateBy { ansi -> ansi.ansiCode } + } + + fun pack(vararg options: ANSI): UShort = options.fold(0u) { acc, ansi -> + acc or (1u shl ansi.bitIdx).toUShort() + } + + fun unpack(packed: UShort): List { + val res = mutableListOf() + var remaining = packed + for(i in 0..15) { + if((remaining and 1u) != 0.toUShort()) { + _mapping[i]?.let { res += it } + } + remaining = (remaining.toUInt() shr 1).toUShort() + } + return res + } + } +} + +@Serializable +sealed class ClientMonitorEvent { + @Serializable + data class Message(val msg: String, val stream: Stream, val options: UShort) : ClientMonitorEvent() { + constructor(msg: String, stream: Stream, vararg options: ANSI) : this(msg, stream, ANSI.pack(*options)) + constructor(msg: String, stream: Stream, options: List) : this(msg, stream, ANSI.pack(*options.toTypedArray())) + } + + @Serializable + data class ApplicationFinished(val exitCode: Int) : ClientMonitorEvent() +} + +@Serializable +sealed class ServerMonitorEvent { + sealed interface ITextEvent + @Serializable + object Cleared : ServerMonitorEvent() + + @Serializable + data class ApplicationStart(val time: Instant) : ServerMonitorEvent(), ITextEvent + + @Serializable + data class Message(val timeStamp: Instant, val message: String, val options: UShort, val stream: Stream) : ServerMonitorEvent(), ITextEvent + + @Serializable + data class Backlog(val start: ApplicationStart?, val messages: List, val end: ApplicationEnd?) : ServerMonitorEvent() + + @Serializable + data class ApplicationEnd(val time: Instant, val exitCode: Int) : ServerMonitorEvent(), ITextEvent } \ No newline at end of file diff --git a/common/src/commonMain/kotlin/com/jaytux/phoebench/common/Routes.kt b/common/src/commonMain/kotlin/com/jaytux/phoebench/common/Routes.kt index 4953632..b0daf13 100644 --- a/common/src/commonMain/kotlin/com/jaytux/phoebench/common/Routes.kt +++ b/common/src/commonMain/kotlin/com/jaytux/phoebench/common/Routes.kt @@ -29,6 +29,7 @@ object Routes { val get = ApiRoute.getUuid("/project", Elevation.AUTH) val update = ApiRoute.patchUuidNoRes("/project", Elevation.AUTH) val delete = ApiRoute.deleteUuidNoRes("/project", Elevation.AUTH) + val rmMonitor = ApiRoute.deleteUuidNoRes("/project/monitor", Elevation.AUTH) } object Benchmark { @@ -52,5 +53,10 @@ object Routes { val home = SSERoute.noArgs("/rt/home", Elevation.AUTH) val admin = SSERoute.noArgs("/rt/admin", Elevation.ADMIN) val projectSpecific = SSERoute.uuid("/rt/project", Elevation.AUTH) + val monitor = SSERoute.uuid("/rt/monitor", Elevation.AUTH) + } + + object CSE { + val monitor = CSERoute.uuid("/stream/monitor", Elevation.AUTH) } } \ No newline at end of file diff --git a/gradle/libs.versions.toml b/gradle/libs.versions.toml index da22431..4a5e18c 100644 --- a/gradle/libs.versions.toml +++ b/gradle/libs.versions.toml @@ -23,6 +23,7 @@ kolor-picker = "2.1.0" shadow = "9.3.0" clikt = "5.0.3" buildconfig = "6.0.10" +process = "1.5.1" [libraries] androidx-lifecycle-viewmodel = { group = "org.jetbrains.androidx.lifecycle", name = "lifecycle-viewmodel", version.ref = "androidx-lifecycle" } @@ -50,6 +51,7 @@ ktor-client-content-negotiation = { module = "io.ktor:ktor-client-content-negoti ktor-serialization-kotlinx-json = { module = "io.ktor:ktor-serialization-kotlinx-json", version.ref = "ktor" } ktor-client-auth = { module = "io.ktor:ktor-client-auth", version.ref = "ktor" } ktor-client-okhttp = { module = "io.ktor:ktor-client-okhttp", version.ref = "ktor" } +ktor-client-websocket = { module = "io.ktor:ktor-client-websockets", version.ref = "ktor" } ktor-server-content-negotiation = { module = "io.ktor:ktor-server-content-negotiation", version.ref = "ktor" } ktor-server-call-logging = { module = "io.ktor:ktor-server-call-logging", version.ref = "ktor" } @@ -64,6 +66,7 @@ ktor-server-auth-jwt = { module = "io.ktor:ktor-server-auth-jwt", version.ref = ktor-server-status-pages = { module = "io.ktor:ktor-server-status-pages", version.ref = "ktor" } ktor-server-cors = { module = "io.ktor:ktor-server-cors", version.ref = "ktor" } ktor-server-sse = { module = "io.ktor:ktor-server-sse", version.ref = "ktor" } +ktor-server-websocket = { module = "io.ktor:ktor-server-websockets", version.ref = "ktor" } json = { module = "org.json:json", version.ref = "json" } kotlinx-datetime = { module = "org.jetbrains.kotlinx:kotlinx-datetime", version.ref = "datetime" } @@ -87,6 +90,7 @@ kolor = { module = "com.kborowy:kolor-picker", version.ref = "kolor-picker" } java-keystore = { module = "com.github.javakeyring:java-keyring", version.ref = "java-keystore" } clikt = { module = "com.github.ajalt.clikt:clikt", version.ref = "clikt" } +process = { module = "com.github.pgreze:kotlin-process", version.ref = "process" } [plugins] composeMultiplatform = { id = "org.jetbrains.compose", version.ref = "compose-multiplatform" } diff --git a/server/build.gradle.kts b/server/build.gradle.kts index a567c34..8e4fae4 100644 --- a/server/build.gradle.kts +++ b/server/build.gradle.kts @@ -7,7 +7,7 @@ plugins { } group = "com.jaytux.phoebench" -version = PhoebenchVersion(1, 1, 2) +version = rootProject.version as PhoebenchVersion if((version as PhoebenchVersion) < (rootProject.version as PhoebenchVersion)) throw GradleException("Server version must be at least as high as protocol/common version") @@ -56,6 +56,7 @@ dependencies { implementation(libs.ktor.server.cors) implementation(libs.ktor.server.status.pages) implementation(libs.ktor.server.sse) + implementation(libs.ktor.server.websocket) implementation(libs.ktor.serialization.kotlinx.json) diff --git a/server/src/main/kotlin/com/jaytux/phoebench/server/Buses.kt b/server/src/main/kotlin/com/jaytux/phoebench/server/Buses.kt index 054b776..cf7e463 100644 --- a/server/src/main/kotlin/com/jaytux/phoebench/server/Buses.kt +++ b/server/src/main/kotlin/com/jaytux/phoebench/server/Buses.kt @@ -3,6 +3,8 @@ package com.jaytux.phoebench.server import com.jaytux.phoebench.common.AdminEvent import com.jaytux.phoebench.common.HomeEvent import com.jaytux.phoebench.common.ProjectEvent +import com.jaytux.phoebench.common.ServerMonitorEvent +import io.ktor.client.plugins.api.MonitoringEvent import io.ktor.util.reflect.typeInfo import kotlinx.serialization.serializer import java.util.concurrent.ConcurrentHashMap @@ -11,6 +13,10 @@ import kotlin.uuid.Uuid object Buses { private val _projectBuses = ConcurrentHashMap>() + private val _monitorBuses = ConcurrentHashMap>() val homeBus = SSEBus(typeOf(), serializer()) val adminBus = SSEBus(typeOf(), serializer()) @@ -19,9 +25,16 @@ object Buses { SSEBus(typeOf(), serializer()) } + fun monitorBus(id: Uuid) = _monitorBuses.computeIfAbsent(id) { + SSEBus.MonitorSSEBus( + SSEBus(typeOf(), serializer()) + ) { start, events, end -> ServerMonitorEvent.Backlog(start, events, end) } + } + fun allBuses(): List> { - val res = ArrayList>(_projectBuses.size + 2) + val res = ArrayList>(_projectBuses.size + _monitorBuses.size + 2) res.addAll(_projectBuses.values) + res.addAll(_monitorBuses.values.map { it.bus }) res.add(homeBus) res.add(adminBus) return res diff --git a/server/src/main/kotlin/com/jaytux/phoebench/server/Main.kt b/server/src/main/kotlin/com/jaytux/phoebench/server/Main.kt index 44ed0df..6306b04 100644 --- a/server/src/main/kotlin/com/jaytux/phoebench/server/Main.kt +++ b/server/src/main/kotlin/com/jaytux/phoebench/server/Main.kt @@ -4,10 +4,15 @@ import com.jaytux.phoebench.common.* import com.jaytux.phoebench.server.db.DB import com.jaytux.phoebench.server.db.Meta import com.jaytux.phoebench.server.db.Metas +import com.jaytux.phoebench.server.db.Project import com.jaytux.phoebench.server.db.User import com.jaytux.phoebench.server.handlers.* import com.jaytux.phoebench.server.handlers.ProjectHandler.accessibleProject +import com.jaytux.phoebench.server.handlers.ProjectHandler.accessibleProjectCSE import io.ktor.http.* +import io.ktor.serialization.WebsocketContentConverter +import io.ktor.serialization.deserialize +import io.ktor.serialization.kotlinx.KotlinxWebsocketSerializationConverter import io.ktor.serialization.kotlinx.json.* import io.ktor.server.application.* import io.ktor.server.auth.* @@ -21,8 +26,14 @@ import io.ktor.server.request.* import io.ktor.server.response.* import io.ktor.server.routing.* import io.ktor.server.sse.* +import io.ktor.server.websocket.WebSockets +import io.ktor.server.websocket.webSocket +import io.ktor.websocket.CloseReason +import io.ktor.websocket.close +import kotlinx.coroutines.channels.ClosedReceiveChannelException import kotlinx.serialization.json.Json import org.jetbrains.exposed.v1.jdbc.transactions.transaction +import kotlin.time.Clock import kotlin.uuid.Uuid fun Application.module() { @@ -78,6 +89,10 @@ fun Application.module() { install(SSE) {} + install(WebSockets) { + contentConverter = KotlinxWebsocketSerializationConverter(Json) + } + install(StatusPages) { status(HttpStatusCode.Forbidden) { call, status -> call.respond(status, ErrorResponse("Access Forbidden: CORS failed.")) @@ -137,6 +152,7 @@ fun Application.module() { getAuth(Routes.Project.get, ProjectHandler::getProject) patchAuth(Routes.Project.update, ProjectHandler::updateProject) deleteAuth(Routes.Project.delete, ProjectHandler::deleteProject) + deleteAuth(Routes.Project.rmMonitor, ProjectHandler::deleteMonitor) postAuth(Routes.Benchmark.new, ProjectHandler::createBenchmark) patchAuth(Routes.Benchmark.update, ProjectHandler::updateBenchmark) @@ -149,13 +165,38 @@ fun Application.module() { postAuth(Routes.Entry.new, ProjectHandler::createEntry) deleteAuth(Routes.Entry.delete, ProjectHandler::deleteEntry) - sseAuth(Routes.SSE.home, { _, _ -> }) { _, _ -> Buses.homeBus } - sseAdmin(Routes.SSE.admin) { _, _ -> Buses.adminBus } + sseAuth(Routes.SSE.home, { _, _ -> }) { _, _, _ -> Buses.homeBus } + sseAdmin(Routes.SSE.admin) { _, _, _ -> Buses.adminBus } sseAuth(Routes.SSE.projectSpecific, { user, uuid -> transaction { accessibleProject(user, uuid, false) } } - ) { _, id -> Buses.projectBus(id) } + ) { _, id, _ -> Buses.projectBus(id) } + sseAuth(Routes.SSE.monitor, + { user, uuid -> + transaction { accessibleProject(user, uuid, true) } + } + ) { _, id, sender -> + val bus = Buses.monitorBus(id) + val backlog = bus.onConnect() + println("Bus [$id] connected; sending backlog (${backlog.start != null} // ${backlog.messages.size} // ${backlog.end != null})") + sender(bus.bus.serializer, backlog) + bus.bus + } + + cseAuth(Routes.CSE.monitor, + { id, user -> transaction { accessibleProjectCSE(user, id, true) } }, + ProjectHandler::monitorSetup) { _, _, bus, event -> + try { + when(event) { + is ClientMonitorEvent.ApplicationFinished -> bus.end(ServerMonitorEvent.ApplicationEnd(Clock.System.now(), event.exitCode)) + is ClientMonitorEvent.Message -> bus.event(ServerMonitorEvent.Message(Clock.System.now(), event.msg, event.options, event.stream)) + } + } + catch(e: IllegalStateException) { + throw RouteError.CSERouteError("Could not deliver message: ${e.message ?: "unknown error"}", CloseReasons.CONFLICT) + } + } } get("{...}") { diff --git a/server/src/main/kotlin/com/jaytux/phoebench/server/MutableBackLog.kt b/server/src/main/kotlin/com/jaytux/phoebench/server/MutableBackLog.kt new file mode 100644 index 0000000..732c20a --- /dev/null +++ b/server/src/main/kotlin/com/jaytux/phoebench/server/MutableBackLog.kt @@ -0,0 +1,38 @@ +package com.jaytux.phoebench.server + +class MutableBackLog { + var start: TStart? = null + private set + private val _events = mutableListOf() + val events = _events.immutable() + var end: TEnd? = null + private set + + fun registerStart(event: TStart) { + if(start != null) throw IllegalStateException("Start is already set.") + start = event + println("[BACKLOG]: start '$event'") + } + + fun onEvent(event: TEvent) { + if(start == null) throw IllegalStateException("Start is not set yet.") + if(end != null) throw IllegalStateException("End is already set.") + _events += event + println("[BACKLOG]: event '$event'") + } + + fun registerEnd(event: TEnd) { + if(start == null) throw IllegalStateException("Start is not set yet.") + if(end != null) throw IllegalStateException("End is already set.") + end = event + println("[BACKLOG]: end '$event'") + } + + fun reset() { + start = null + _events.clear() + end = null + } + + fun isRunning() = start != null && end == null +} \ No newline at end of file diff --git a/server/src/main/kotlin/com/jaytux/phoebench/server/SSEBus.kt b/server/src/main/kotlin/com/jaytux/phoebench/server/SSEBus.kt index d0819ec..e826c1a 100644 --- a/server/src/main/kotlin/com/jaytux/phoebench/server/SSEBus.kt +++ b/server/src/main/kotlin/com/jaytux/phoebench/server/SSEBus.kt @@ -61,4 +61,32 @@ class SSEBus(private val _containedType: KType, val serializer: KSerializer( + val bus: SSEBus, val backlog: MutableBackLog = MutableBackLog(), + val mkBacklog: (start: TStart?, events: List, end: TEnd?) -> TBacklog + ) { + val start: TStart? get() = backlog.start + val events: List get() = backlog.events + val end: TEnd? get() = backlog.end + + suspend fun start(start: TStart) { + backlog.registerStart(start) + bus.sendAll(start) + } + + suspend fun event(event: TEvent) { + backlog.onEvent(event) + bus.sendAll(event) + } + + suspend fun end(end: TEnd) { + backlog.registerEnd(end) + bus.sendAll(end) + } + + suspend fun onConnect() = mkBacklog(start, events, end) + + fun isRunning() = backlog.isRunning() + } } \ No newline at end of file diff --git a/server/src/main/kotlin/com/jaytux/phoebench/server/Util.kt b/server/src/main/kotlin/com/jaytux/phoebench/server/Util.kt index 2c37b18..9cb0cf3 100644 --- a/server/src/main/kotlin/com/jaytux/phoebench/server/Util.kt +++ b/server/src/main/kotlin/com/jaytux/phoebench/server/Util.kt @@ -23,4 +23,6 @@ fun nowPlus(time: Int, unit: DateTimeUnit): Instant { return now.plus(time, unit, systemTZ) } -infix fun Pair.app(t3: T3) = Triple(first, second, t3) \ No newline at end of file +infix fun Pair.app(t3: T3) = Triple(first, second, t3) + +fun MutableList.immutable(): List = this \ No newline at end of file diff --git a/server/src/main/kotlin/com/jaytux/phoebench/server/handlers/Bridge.kt b/server/src/main/kotlin/com/jaytux/phoebench/server/handlers/Bridge.kt index 668b67e..b24a7c9 100644 --- a/server/src/main/kotlin/com/jaytux/phoebench/server/handlers/Bridge.kt +++ b/server/src/main/kotlin/com/jaytux/phoebench/server/handlers/Bridge.kt @@ -4,18 +4,19 @@ import com.jaytux.phoebench.server.handlers.RouteError.Companion.wrapped import com.jaytux.phoebench.server.handlers.RouteError.Companion.wrappedAdmin import com.jaytux.phoebench.server.handlers.RouteError.Companion.wrappedAuth import com.jaytux.phoebench.common.ApiRoute -import com.jaytux.phoebench.common.Either +import com.jaytux.phoebench.common.CSERoute +import com.jaytux.phoebench.common.CloseReasons import com.jaytux.phoebench.common.Elevation import com.jaytux.phoebench.common.EmptyRequest import com.jaytux.phoebench.common.ErrorResponse import com.jaytux.phoebench.common.SSERoute import com.jaytux.phoebench.common.foldSuspend -import com.jaytux.phoebench.server.Auth.setup import com.jaytux.phoebench.server.SSEBus import com.jaytux.phoebench.server.db.User -import com.jaytux.phoebench.server.handlers.RouteError import io.ktor.http.ContentType import io.ktor.http.HttpStatusCode +import io.ktor.serialization.deserialize +import io.ktor.serialization.kotlinx.KotlinxWebsocketSerializationConverter import io.ktor.server.application.ApplicationCall import io.ktor.server.auth.jwt.JWTPrincipal import io.ktor.server.auth.principal @@ -31,10 +32,15 @@ import io.ktor.server.routing.delete import io.ktor.server.routing.patch import io.ktor.server.sse.heartbeat import io.ktor.server.sse.sse +import io.ktor.server.websocket.webSocket import io.ktor.sse.ServerSentEvent import io.ktor.util.reflect.typeInfo import io.ktor.utils.io.CancellationException +import io.ktor.websocket.CloseReason +import io.ktor.websocket.close +import kotlinx.coroutines.channels.ClosedReceiveChannelException import kotlinx.coroutines.flow.SharedFlow +import kotlinx.serialization.KSerializer import kotlinx.serialization.json.Json import org.jetbrains.exposed.v1.jdbc.transactions.transaction import kotlin.time.Duration.Companion.seconds @@ -165,7 +171,7 @@ inline fun Route.patchAdmin(api: ApiRoute inline fun Route.wrapSSE( api: SSERoute, noinline extra: suspend (ApplicationCall, TParams) -> TInter, - noinline prepare: suspend (TInter, TParams) -> SSEBus, + noinline prepare: suspend (TInter, TParams, sender: suspend (KSerializer, TEvent) -> Unit) -> SSEBus, noinline extract: suspend (SSEBus, TInter, TParams) -> SharedFlow, noinline handler: suspend (TFlow, sender: suspend (TEvent) -> Unit) -> Unit, noinline onCancel: suspend (SSEBus, TInter, TParams, CancellationException) -> Unit @@ -182,11 +188,14 @@ inline fun Route.wrap ) val inter = extra(call, params) - val bus = prepare(inter, params) + val bus = prepare(inter, params) { serializer, event -> send(ServerSentEvent(data = Json.encodeToString(serializer, event))) } val stream = extract(bus, inter, params) try { stream.collect { - handler(it) { ev -> send(ServerSentEvent(data = Json.encodeToString(bus.serializer, ev))) } + handler(it) { ev -> + println("[SSE ${api.pattern}]: sending event!") + send(ServerSentEvent(data = Json.encodeToString(bus.serializer, ev))) + } } } catch(e: CancellationException) { @@ -214,7 +223,7 @@ inline fun Route.wrap inline fun Route.wrapAuthSSE( api: SSERoute, noinline verifyUser: suspend (User, TParams) -> Unit, - noinline prepare: suspend (User, TParams) -> SSEBus + noinline prepare: suspend (User, TParams, sender: suspend (KSerializer, TEvent) -> Unit) -> SSEBus ) = wrapSSE(api, extra = { call, params -> val principal = call.principal() @@ -239,23 +248,113 @@ inline fun Route.wrapAuthSSE( onCancel = { bus, user, _, _ -> bus.disconnect(user.id.value) } ) -inline fun Route.sse(api: SSERoute, noinline setup: suspend (TParams) -> SSEBus) = - wrapSSE(api, +inline fun Route.sse(api: SSERoute, noinline prepare: suspend (TParams, sender: suspend (KSerializer, TEvent) -> Unit) -> SSEBus) { + if(api.elevation != Elevation.UN_AUTH) throw IllegalArgumentException("SSE ${api.pattern} can only be used with ${api.elevation}") + wrapSSE( + api, extra = { _, _ -> }, - prepare = { _, params -> setup(params) }, + prepare = { _, params, sender -> prepare(params, sender) }, extract = { bus, _, _ -> bus.unRegistered() }, handler = { it, sender -> sender(it) }, onCancel = { _, _, _, _ -> } ) +} inline fun Route.sseAuth(api: SSERoute, - noinline verifyUser: suspend (User, TParams) -> Unit, noinline setup: suspend (User, TParams) -> SSEBus -) = wrapAuthSSE(api, verifyUser, setup) + noinline verifyUser: suspend (User, TParams) -> Unit, noinline prepare: suspend (User, TParams, sender: suspend (KSerializer, TEvent) -> Unit) -> SSEBus +) { + if(api.elevation != Elevation.AUTH) throw IllegalArgumentException("SSE ${api.pattern} can only be used with ${api.elevation}") + wrapAuthSSE(api, verifyUser, prepare) +} inline fun Route.sseAdmin(api: SSERoute, - noinline setup: suspend (User, TParams) -> SSEBus -) = wrapAuthSSE(api, { user, _ -> - if(!user.isAdmin) { - throw RouteError("Admin access required", HttpStatusCode.Forbidden) + noinline prepare: suspend (User, TParams, sender: suspend (KSerializer, TEvent) -> Unit) -> SSEBus +) { + if(api.elevation != Elevation.ADMIN) throw IllegalArgumentException("SSE ${api.pattern} can only be used with ${api.elevation}") + wrapAuthSSE(api, { user, _ -> + if(!user.isAdmin) { + throw RouteError("Admin access required", HttpStatusCode.Forbidden) + } + }, prepare) +} + +inline fun Route.wrapCSE( + api: CSERoute, + noinline extra: suspend (ApplicationCall, TParams) -> TInter, + noinline setup: suspend (TParams, TInter) -> TExtra, + noinline handler: suspend (TParams, TInter, TExtra, TEvent) -> Unit +) { + webSocket(api.pattern) { + try { + val params = api.extractParams(call.parameters) ?: + throw RouteError.CSERouteError("Missing or malformed parameters for CSE ${api.pattern}", CloseReasons.INVALID_REQUEST) + val inter = extra(call, params) + + val converter = KotlinxWebsocketSerializationConverter(Json) + val extra = setup(params, inter) + for(frame in incoming) { + handler(params, inter, extra, converter.deserialize(frame)) + } + } + catch(e: ClosedReceiveChannelException) {} + catch(e: RouteError.CSERouteError) { + close(CloseReason(e.code, e.message ?: "Unknown error")) + } + catch(e: Exception) { + close(CloseReason(CloseReason.Codes.INTERNAL_ERROR, e.message ?: "Unknown error")) + } } -}, setup) \ No newline at end of file +} + +inline fun Route.wrapAuthCSE( + api: CSERoute, + noinline verifyUser: suspend (TParams, User) -> Unit, + noinline setup: suspend (TParams, User) -> TExtra, + noinline handler: suspend (TParams, User, TExtra, TEvent) -> Unit +) = wrapCSE( + api = api, + extra = { call, params -> + val principal = call.principal() + val userId = principal?.payload?.getClaim(com.jaytux.phoebench.common.Auth.JWT_CLAIM)?.asString() + ?: throw RouteError.CSERouteError("Missing user claim", CloseReasons.NOT_AUTHORIZED) + val user = transaction { + User.findById(Uuid.parse(userId)) ?: throw RouteError.CSERouteError( + "User not found", + CloseReasons.NOT_AUTHORIZED + ) + } + verifyUser(params, user) + user + }, + setup = setup, + handler = handler +) + +inline fun Route.cse( + api: CSERoute, + noinline setup: suspend (TParams, Unit) -> TExtra, + noinline handler: suspend (TParams, TExtra, TEvent) -> Unit +) { + if(api.elevation != Elevation.UN_AUTH) throw IllegalArgumentException("CSE ${api.pattern} can only be used with ${api.elevation}") + wrapCSE(api, { _, _ -> }, setup) { params, _, extra, event -> handler(params, extra, event) } +} + +inline fun Route.cseAuth( + api: CSERoute, noinline verifyUser: suspend (TParams, User) -> Unit, + noinline setup: suspend (TParams, User) -> TExtra, noinline handler: suspend (TParams, User, TExtra, TEvent) -> Unit +) { + if(api.elevation != Elevation.AUTH) throw IllegalArgumentException("CSE ${api.pattern} can only be used with ${api.elevation}") + wrapAuthCSE(api, verifyUser, setup, handler) +} + +inline fun Route.cseAdmin( + api: CSERoute, noinline setup: suspend (TParams, User) -> TExtra, + noinline handler: suspend (TParams, User, TExtra, TEvent) -> Unit +) { + if(api.elevation != Elevation.ADMIN) throw IllegalArgumentException("CSE ${api.pattern} can only be used with ${api.elevation}") + wrapAuthCSE(api, { _, user -> + if(!user.isAdmin) { + throw RouteError.CSERouteError("Admin access required", CloseReasons.NOT_AUTHORIZED) + } + }, setup, handler) +} \ No newline at end of file diff --git a/server/src/main/kotlin/com/jaytux/phoebench/server/handlers/ProjectHandler.kt b/server/src/main/kotlin/com/jaytux/phoebench/server/handlers/ProjectHandler.kt index e470e81..dcbac2d 100644 --- a/server/src/main/kotlin/com/jaytux/phoebench/server/handlers/ProjectHandler.kt +++ b/server/src/main/kotlin/com/jaytux/phoebench/server/handlers/ProjectHandler.kt @@ -3,6 +3,8 @@ package com.jaytux.phoebench.server.handlers import com.jaytux.phoebench.common.BenchmarkRequest import com.jaytux.phoebench.common.BenchmarkResponse import com.jaytux.phoebench.common.BenchmarkSummary +import com.jaytux.phoebench.common.ClientMonitorEvent +import com.jaytux.phoebench.common.CloseReasons import com.jaytux.phoebench.common.EmptyRequest import com.jaytux.phoebench.common.EmptyResponse import com.jaytux.phoebench.common.EntryRequest @@ -19,7 +21,9 @@ import com.jaytux.phoebench.common.PartialVersionRequest import com.jaytux.phoebench.common.ProjectEvent import com.jaytux.phoebench.common.ProjectRequest import com.jaytux.phoebench.common.ProjectResponse +import com.jaytux.phoebench.common.ServerMonitorEvent import com.jaytux.phoebench.server.Buses +import com.jaytux.phoebench.server.SSEBus import com.jaytux.phoebench.server.app import com.jaytux.phoebench.server.db.Benchmark import com.jaytux.phoebench.server.db.Benchmarks @@ -43,6 +47,7 @@ import org.jetbrains.exposed.v1.core.inList import org.jetbrains.exposed.v1.core.notInList import org.jetbrains.exposed.v1.jdbc.deleteWhere import org.jetbrains.exposed.v1.jdbc.transactions.transaction +import kotlin.time.Clock import kotlin.uuid.Uuid object ProjectHandler { @@ -68,6 +73,11 @@ object ProjectHandler { return project.isAccessible(user, forEditing) } + fun Transaction.accessibleProjectCSE(user: User, id: Uuid, forEditing: Boolean): Project { + val project = Project.findById(id) ?: throw RouteError.CSERouteError("Invalid project ID.", CloseReasons.NOT_FOUND) + return project.isAccessible(user, forEditing) + } + fun Transaction.accessibleBenchmark(user: User, id: Uuid, forEditing: Boolean): Pair { val bench = Benchmark.findById(id) ?: throw RouteError("Invalid benchmark ID.", HttpStatusCode.NotFound) return bench.isAccessible(user, forEditing) @@ -381,4 +391,26 @@ object ProjectHandler { return home(user, EmptyRequest()) } + + suspend fun monitorSetup(id: Uuid, user: User) = Buses.monitorBus(id).also { + try { + if(!it.isRunning()) it.backlog.reset() + it.start(ServerMonitorEvent.ApplicationStart(Clock.System.now())) + } + catch(e: IllegalStateException) { + throw RouteError.CSERouteError("Monitor may already be running. Please refrain from starting a second trace.", CloseReasons.CONFLICT) + } + } + + fun deleteMonitor(user: User, req: Uuid) = transaction { + checkMigration(user) + accessibleProject(user, req, true) + val bus = Buses.monitorBus(req) + if(bus.isRunning()) throw RouteError("Monitor is still running", HttpStatusCode.Conflict) + bus.backlog.reset() + ServerScope.launch { + bus.bus.sendAll(ServerMonitorEvent.Cleared) + } + success(EmptyResponse()) + } } \ No newline at end of file diff --git a/server/src/main/kotlin/com/jaytux/phoebench/server/handlers/RouteError.kt b/server/src/main/kotlin/com/jaytux/phoebench/server/handlers/RouteError.kt index f3510b5..26263e0 100644 --- a/server/src/main/kotlin/com/jaytux/phoebench/server/handlers/RouteError.kt +++ b/server/src/main/kotlin/com/jaytux/phoebench/server/handlers/RouteError.kt @@ -1,6 +1,7 @@ package com.jaytux.phoebench.server.handlers import com.jaytux.phoebench.common.Auth +import com.jaytux.phoebench.common.CloseReasons import com.jaytux.phoebench.common.ErrorResponse import com.jaytux.phoebench.server.db.User import io.ktor.http.ContentType @@ -13,6 +14,7 @@ import io.ktor.server.response.respondText import io.ktor.server.routing.RoutingCall import io.ktor.server.routing.RoutingContext import io.ktor.util.logging.KtorSimpleLogger +import io.ktor.websocket.CloseReason import kotlinx.serialization.json.Json import org.jetbrains.exposed.v1.jdbc.transactions.transaction import kotlin.uuid.Uuid @@ -75,4 +77,9 @@ open class RouteError(message: String, val status: HttpStatusCode = HttpStatusCo fun unauthorized(message: String): Nothing = throw RouteError(message, HttpStatusCode.Unauthorized) } + + class CSERouteError(message: String, val code: Short) : RouteError(message, HttpStatusCode.NotImplemented) { + constructor(message: String, code: CloseReason.Codes) : this(message, code.code) + constructor(message: String, code: CloseReasons) : this(message, code.code) + } } \ No newline at end of file