From 3bf9d834509bf25460d394067ed740a7dfb4df65 Mon Sep 17 00:00:00 2001 From: observer Date: Sat, 12 Sep 2026 20:36:24 +0100 Subject: [PATCH] Use shared flow (buffer) in reader --- .../shr4pnel/ferretirc/net/MessageReader.kt | 20 +++++++++++-------- 1 file changed, 12 insertions(+), 8 deletions(-) diff --git a/src/main/kotlin/com/shr4pnel/ferretirc/net/MessageReader.kt b/src/main/kotlin/com/shr4pnel/ferretirc/net/MessageReader.kt index 83e7bf9..4f98f2f 100644 --- a/src/main/kotlin/com/shr4pnel/ferretirc/net/MessageReader.kt +++ b/src/main/kotlin/com/shr4pnel/ferretirc/net/MessageReader.kt @@ -11,6 +11,7 @@ import io.ktor.utils.io.ByteReadChannel import io.ktor.utils.io.readLineStrict import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.flow.MutableSharedFlow +import kotlinx.coroutines.flow.asSharedFlow import kotlinx.coroutines.flow.filter import kotlinx.coroutines.flow.first import kotlinx.coroutines.flow.map @@ -21,8 +22,8 @@ import kotlin.reflect.cast class MessageReader(socket: Socket, scope: CoroutineScope) : MessageIO(socket, scope) { private lateinit var receive: ByteReadChannel - private val listenedMessages = mutableListOf() - val messageBuffer = MutableSharedFlow(16, 64) + private val messageBuffer = MutableSharedFlow(16, 64) + val sharedMessageBuffer = messageBuffer.asSharedFlow() val parser = MessageParser(incomingMessages) private val logger = if (Helpers.enableLogging) KotlinLogging.logger {} else KotlinLogging.logger(NOPLogger.NOP_LOGGER) @@ -32,16 +33,19 @@ class MessageReader(socket: Socket, scope: CoroutineScope) : MessageIO(socket, s parser.start() } scope.launch { - for (msg in parser.incomingParsedMessages) messageBuffer.emit(msg) + for (msg in parser.incomingParsedMessages) + messageBuffer.emit(msg) } } - override fun start() = scope.launch { - receive = socket.openReadChannel() - while (true) { - val line = receive.readLineStrict() ?: break // >:( no LineEnding option for just CRLF? charlatans... - incomingMessages.send(line) + override fun start() { + scope.launch { + receive = socket.openReadChannel() + while (true) { + val line = receive.readLineStrict() ?: break // >:( no LineEnding option for just CRLF? charlatans... logger.debug { "Receive: $line" } + incomingMessages.send(line) + } } }