Script monitoring

This commit is contained in:
2026-09-01 11:47:10 +02:00
parent 8da391b6ce
commit 85902620da
37 changed files with 803 additions and 71 deletions
@@ -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<Uuid, SSEBus<ProjectEvent>>()
private val _monitorBuses = ConcurrentHashMap<Uuid, SSEBus.MonitorSSEBus<
ServerMonitorEvent, ServerMonitorEvent.ApplicationStart, ServerMonitorEvent.Message,
ServerMonitorEvent.ApplicationEnd, ServerMonitorEvent.Backlog
>>()
val homeBus = SSEBus<HomeEvent>(typeOf<HomeEvent>(), serializer<HomeEvent>())
val adminBus = SSEBus<AdminEvent>(typeOf<AdminEvent>(), serializer<AdminEvent>())
@@ -19,9 +25,16 @@ object Buses {
SSEBus<ProjectEvent>(typeOf<ProjectEvent>(), serializer<ProjectEvent>())
}
fun monitorBus(id: Uuid) = _monitorBuses.computeIfAbsent(id) {
SSEBus.MonitorSSEBus(
SSEBus(typeOf<ServerMonitorEvent>(), serializer<ServerMonitorEvent>())
) { start, events, end -> ServerMonitorEvent.Backlog(start, events, end) }
}
fun allBuses(): List<SSEBus<*>> {
val res = ArrayList<SSEBus<*>>(_projectBuses.size + 2)
val res = ArrayList<SSEBus<*>>(_projectBuses.size + _monitorBuses.size + 2)
res.addAll(_projectBuses.values)
res.addAll(_monitorBuses.values.map { it.bus })
res.add(homeBus)
res.add(adminBus)
return res
@@ -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("{...}") {
@@ -0,0 +1,38 @@
package com.jaytux.phoebench.server
class MutableBackLog<TStart, TEvent, TEnd> {
var start: TStart? = null
private set
private val _events = mutableListOf<TEvent>()
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
}
@@ -61,4 +61,32 @@ class SSEBus<T>(private val _containedType: KType, val serializer: KSerializer<T
map
}.forEach { it.second.emit(Cancellation.error()) }
}
class MonitorSSEBus<TSuper, TStart : TSuper, TEvent : TSuper, TEnd: TSuper, TBacklog: TSuper>(
val bus: SSEBus<TSuper>, val backlog: MutableBackLog<TStart, TEvent, TEnd> = MutableBackLog(),
val mkBacklog: (start: TStart?, events: List<TEvent>, end: TEnd?) -> TBacklog
) {
val start: TStart? get() = backlog.start
val events: List<TEvent> 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()
}
}
@@ -23,4 +23,6 @@ fun nowPlus(time: Int, unit: DateTimeUnit): Instant {
return now.plus(time, unit, systemTZ)
}
infix fun <T1, T2, T3> Pair<T1, T2>.app(t3: T3) = Triple(first, second, t3)
infix fun <T1, T2, T3> Pair<T1, T2>.app(t3: T3) = Triple(first, second, t3)
fun <T> MutableList<T>.immutable(): List<T> = this
@@ -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 <reified TReq: Any, reified TRes: Any> Route.patchAdmin(api: ApiRoute
inline fun <reified TParams: Any, reified TEvent: Any, TInter, TFlow> Route.wrapSSE(
api: SSERoute<TParams, TEvent>, noinline extra: suspend (ApplicationCall, TParams) -> TInter,
noinline prepare: suspend (TInter, TParams) -> SSEBus<TEvent>,
noinline prepare: suspend (TInter, TParams, sender: suspend (KSerializer<TEvent>, TEvent) -> Unit) -> SSEBus<TEvent>,
noinline extract: suspend (SSEBus<TEvent>, TInter, TParams) -> SharedFlow<TFlow>,
noinline handler: suspend (TFlow, sender: suspend (TEvent) -> Unit) -> Unit,
noinline onCancel: suspend (SSEBus<TEvent>, TInter, TParams, CancellationException) -> Unit
@@ -182,11 +188,14 @@ inline fun <reified TParams: Any, reified TEvent: Any, TInter, TFlow> 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 <reified TParams: Any, reified TEvent: Any, TInter, TFlow> Route.wrap
inline fun <reified TParams: Any, reified TEvent: Any> Route.wrapAuthSSE(
api: SSERoute<TParams, TEvent>,
noinline verifyUser: suspend (User, TParams) -> Unit,
noinline prepare: suspend (User, TParams) -> SSEBus<TEvent>
noinline prepare: suspend (User, TParams, sender: suspend (KSerializer<TEvent>, TEvent) -> Unit) -> SSEBus<TEvent>
) = wrapSSE(api,
extra = { call, params ->
val principal = call.principal<JWTPrincipal>()
@@ -239,23 +248,113 @@ inline fun <reified TParams: Any, reified TEvent: Any> Route.wrapAuthSSE(
onCancel = { bus, user, _, _ -> bus.disconnect(user.id.value) }
)
inline fun <reified TParams: Any, reified TEvent: Any> Route.sse(api: SSERoute<TParams, TEvent>, noinline setup: suspend (TParams) -> SSEBus<TEvent>) =
wrapSSE(api,
inline fun <reified TParams: Any, reified TEvent: Any> Route.sse(api: SSERoute<TParams, TEvent>, noinline prepare: suspend (TParams, sender: suspend (KSerializer<TEvent>, TEvent) -> Unit) -> SSEBus<TEvent>) {
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 <reified TParams: Any, reified TEvent: Any> Route.sseAuth(api: SSERoute<TParams, TEvent>,
noinline verifyUser: suspend (User, TParams) -> Unit, noinline setup: suspend (User, TParams) -> SSEBus<TEvent>
) = wrapAuthSSE(api, verifyUser, setup)
noinline verifyUser: suspend (User, TParams) -> Unit, noinline prepare: suspend (User, TParams, sender: suspend (KSerializer<TEvent>, TEvent) -> Unit) -> SSEBus<TEvent>
) {
if(api.elevation != Elevation.AUTH) throw IllegalArgumentException("SSE ${api.pattern} can only be used with ${api.elevation}")
wrapAuthSSE(api, verifyUser, prepare)
}
inline fun <reified TParams: Any, reified TEvent: Any> Route.sseAdmin(api: SSERoute<TParams, TEvent>,
noinline setup: suspend (User, TParams) -> SSEBus<TEvent>
) = wrapAuthSSE(api, { user, _ ->
if(!user.isAdmin) {
throw RouteError("Admin access required", HttpStatusCode.Forbidden)
noinline prepare: suspend (User, TParams, sender: suspend (KSerializer<TEvent>, TEvent) -> Unit) -> SSEBus<TEvent>
) {
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 <reified TParams: Any, reified TEvent: Any, TInter, TExtra> Route.wrapCSE(
api: CSERoute<TParams, TEvent>,
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<TEvent>(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)
}
inline fun <reified TParams: Any, reified TEvent: Any, TExtra> Route.wrapAuthCSE(
api: CSERoute<TParams, TEvent>,
noinline verifyUser: suspend (TParams, User) -> Unit,
noinline setup: suspend (TParams, User) -> TExtra,
noinline handler: suspend (TParams, User, TExtra, TEvent) -> Unit
) = wrapCSE<TParams, TEvent, User, TExtra>(
api = api,
extra = { call, params ->
val principal = call.principal<JWTPrincipal>()
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 <reified TParams: Any, reified TEvent: Any, TExtra> Route.cse(
api: CSERoute<TParams, TEvent>,
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 <reified TParams: Any, reified TEvent: Any, TExtra> Route.cseAuth(
api: CSERoute<TParams, TEvent>, 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 <reified TParams: Any, reified TEvent: Any, TExtra> Route.cseAdmin(
api: CSERoute<TParams, TEvent>, 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)
}
@@ -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<Project, Benchmark> {
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())
}
}
@@ -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)
}
}