diff --git a/fleet/lsp.protocol/BUILD.bazel b/fleet/lsp.protocol/BUILD.bazel index 861b05d932c3..f00e6604791c 100644 --- a/fleet/lsp.protocol/BUILD.bazel +++ b/fleet/lsp.protocol/BUILD.bazel @@ -27,7 +27,6 @@ jvm_library( "@lib//:kotlinx-serialization-json", "//fleet/util/core", "@lib//:jetbrains-annotations", - "//fleet/ktor/network/tls", "@lib//:kotlinx-io-core", ] ) diff --git a/fleet/lsp.protocol/fleet.lsp.protocol.iml b/fleet/lsp.protocol/fleet.lsp.protocol.iml index 41fad3e4396a..96ce2563ba5e 100644 --- a/fleet/lsp.protocol/fleet.lsp.protocol.iml +++ b/fleet/lsp.protocol/fleet.lsp.protocol.iml @@ -36,7 +36,6 @@ - \ No newline at end of file diff --git a/fleet/lsp.protocol/gradlebuild/build.gradle.kts b/fleet/lsp.protocol/gradlebuild/build.gradle.kts index 828a97e04d46..64f9afcb7ced 100644 --- a/fleet/lsp.protocol/gradlebuild/build.gradle.kts +++ b/fleet/lsp.protocol/gradlebuild/build.gradle.kts @@ -70,7 +70,6 @@ kotlin { exclude(group = "org.jetbrains.kotlin", module = "kotlin-stdlib") } implementation(project(":fleet.util.core")) - implementation(project(":fleet.ktor.network.tls")) } // KOTLIN__MARKER_END } \ No newline at end of file diff --git a/fleet/lsp.protocol/srcCommonMain/com/jetbrains/lsp/implementation/io.kt b/fleet/lsp.protocol/srcCommonMain/com/jetbrains/lsp/implementation/io.kt index 7c1d49a85612..045ce6604964 100644 --- a/fleet/lsp.protocol/srcCommonMain/com/jetbrains/lsp/implementation/io.kt +++ b/fleet/lsp.protocol/srcCommonMain/com/jetbrains/lsp/implementation/io.kt @@ -40,7 +40,7 @@ suspend fun ByteReader.readUTF8Line(): String? { if (builder.isNotEmpty() && builder[builder.length - 1] == '\r') { builder.deleteAt(builder.length - 1) } - assert(readBuffer.readByte() == 0x0A.toByte()) + check(readBuffer.readByte() == 0x0A.toByte()) { "expected to see the previously found line terminator" } return builder.toString() } diff --git a/fleet/lsp.protocol/srcCommonMain/com/jetbrains/lsp/protocol/LSP.kt b/fleet/lsp.protocol/srcCommonMain/com/jetbrains/lsp/protocol/LSP.kt index 498f5d96d108..8502c41530f8 100644 --- a/fleet/lsp.protocol/srcCommonMain/com/jetbrains/lsp/protocol/LSP.kt +++ b/fleet/lsp.protocol/srcCommonMain/com/jetbrains/lsp/protocol/LSP.kt @@ -1,7 +1,6 @@ package com.jetbrains.lsp.protocol import fleet.util.isValidUriString -import io.ktor.http.Url import kotlinx.serialization.* import kotlinx.serialization.builtins.serializer import kotlinx.serialization.descriptors.PolymorphicKind @@ -27,26 +26,6 @@ value class URI(val uri: String) { require(uri.isValidUriString()) { "Invalid URI: $uri" } } - /** - * Returns the URI's schema without schema delimiter (`://`) - */ - val scheme: String get() = Url(uri).protocol.name - - /** - * Returns the file name - */ - val fileName: String get() = Url(uri).segments.last() - - /** - * Returns the file extension (without dot) if present - */ - val fileExtension: String? - get() { - val name = fileName - val dotIndex = name.lastIndexOf('.') - return if (dotIndex > 0) name.substring(dotIndex + 1) else null - } - object Schemas { const val FILE: String = "file" const val JRT: String = "jrt" diff --git a/fleet/lsp.protocol/srcJvmMain/com/jetbrains/lsp/implementation/protocolFramingMain.kt b/fleet/lsp.protocol/srcJvmMain/com/jetbrains/lsp/implementation/protocolFramingMain.kt deleted file mode 100644 index 22fcde80cbe5..000000000000 --- a/fleet/lsp.protocol/srcJvmMain/com/jetbrains/lsp/implementation/protocolFramingMain.kt +++ /dev/null @@ -1,37 +0,0 @@ -package com.jetbrains.lsp.implementation - -import com.jetbrains.lsp.protocol.* -import kotlinx.coroutines.Dispatchers -import kotlinx.coroutines.awaitCancellation -import kotlinx.coroutines.runBlocking - -fun main() { - val handler = lspHandlers { - request(Initialize) { initParams -> - InitializeResult( - capabilities = ServerCapabilities( - textDocumentSync = TextDocumentSyncKind.Incremental, - ), - serverInfo = InitializeResult.ServerInfo( - name = "IntelliJ Analyzer", - version = "1.0" - ), - ) - } - notification(DocumentSync.DidOpen) { didOpen -> - println("didOpen: $didOpen") - } - notification(DocumentSync.DidChange) { didChange -> - println("didChange: $didChange") - } - } - runBlocking(Dispatchers.Default) { - tcpServer(TcpConnectionConfig.Server("127.0.0.1", 9999, isMultiClient = true)) { connection -> - withBaseProtocolFraming(connection, exitSignal = null) { incoming, outgoing -> - withLsp(incoming, outgoing, handler) { lsp -> - awaitCancellation() - } - } - } - } -} \ No newline at end of file diff --git a/fleet/lsp.protocol/srcJvmMain/com/jetbrains/lsp/implementation/tcp.kt b/fleet/lsp.protocol/srcJvmMain/com/jetbrains/lsp/implementation/tcp.kt deleted file mode 100644 index d2bde7e2c1e9..000000000000 --- a/fleet/lsp.protocol/srcJvmMain/com/jetbrains/lsp/implementation/tcp.kt +++ /dev/null @@ -1,145 +0,0 @@ -package com.jetbrains.lsp.implementation - -import fleet.util.logging.logger -import io.ktor.network.selector.SelectorManager -import io.ktor.network.sockets.* -import io.ktor.utils.io.ByteReadChannel -import io.ktor.utils.io.ByteWriteChannel -import io.ktor.utils.io.InternalAPI -import kotlinx.coroutines.* -import kotlinx.io.Sink -import kotlinx.io.Source -import kotlin.time.Duration -import kotlin.time.Duration.Companion.seconds - -suspend fun tcpServer(config: TcpConnectionConfig.Server, server: suspend CoroutineScope.(LspConnection) -> Unit) { - SelectorManager(Dispatchers.IO).use { selectorManager -> - aSocket(selectorManager).tcp().bind(config.host, config.port).use { serverSocket -> - LOG.info("Server is listening on ${serverSocket.localAddress}") - - supervisorScope { - var hadClient = false - fun shouldAccept() = !hadClient || config.isMultiClient - - while (shouldAccept()) { - val client = serverSocket.accept() - val clientAddress = client.remoteAddress - hadClient = true - LOG.info("A new client connected at ${clientAddress}") - launch(start = CoroutineStart.ATOMIC) { - try { - client.use { clientSocket -> - coroutineScope { server(KtorSocketConnection(clientSocket)) } - } - } - finally { - LOG.info("Client disconnected ${clientAddress}") - } - } - } - } - } - } -} - - -suspend fun tcpClient( - config: TcpConnectionConfig.Client, - connectionTimeout: Duration = 30.seconds, - body: suspend CoroutineScope.(LspConnection) -> Unit -) { - SelectorManager(Dispatchers.IO).use { selectorManager -> - var backoff = 2.seconds - var timeLeft = connectionTimeout - while (true) { - try { - aSocket(selectorManager).tcp().connect(config.host, config.port).use { server -> - LOG.info("Client is connected to server ${server.remoteAddress}") - try { - coroutineScope { body(KtorSocketConnection(server)) } - } - finally { - LOG.info("Client disconnected from the server") - } - } - return@use - } - catch (e: CancellationException) { - throw e - } - catch (e: Exception) { - if (timeLeft <= Duration.ZERO) throw e - LOG.warn { - "Reconnecting to ${config.host}:${config.port} in $backoff... (error: ${e.message})" - } - val delayTime = backoff.coerceAtMost(timeLeft) - delay(delayTime) - timeLeft -= delayTime - backoff = (backoff * 2).coerceAtMost(connectionTimeout / 2) - } - } - } -} - -sealed interface TcpConnectionConfig { - val host: String - val port: Int - - val isMultiClient: Boolean - - data class Client( - override val host: String, - override val port: Int, - ) : TcpConnectionConfig { - override val isMultiClient: Boolean = false - } - - data class Server( - override val host: String, - override val port: Int, - override val isMultiClient: Boolean, - ) : TcpConnectionConfig -} - -@OptIn(InternalAPI::class) -class KtorByteReader(val input: ByteReadChannel) : ByteReader { - override val closedCause: Throwable? - get() = input.closedCause - override val isClosedForRead: Boolean - get() = input.isClosedForRead - override val readBuffer: Source - get() = input.readBuffer - - override suspend fun awaitContent(min: Int): Boolean = input.awaitContent() - override fun cancel(cause: Throwable?): Unit = input.cancel(cause) -} - -@OptIn(InternalAPI::class) -class KtorByteWriter(val output: ByteWriteChannel) : ByteWriter { - override val isClosedForWrite: Boolean - get() = output.isClosedForWrite - override val closedCause: Throwable? - get() = output.closedCause - override val writeBuffer: Sink - get() = output.writeBuffer - - override suspend fun flush(): Unit = output.flush() - override suspend fun flushAndClose(): Unit = output.flushAndClose() - override fun cancel(cause: Throwable?): Unit = output.cancel(cause) -} - -class KtorSocketConnection(private val socket: Socket) : LspConnection { - override val input: ByteReader = KtorByteReader(socket.openReadChannel()) - override val output: ByteWriter = KtorByteWriter(socket.openWriteChannel(autoFlush = true)) - - override fun close() { - socket.close() - } - - override fun isAlive(): Boolean { - return !socket.isClosed - } -} - - -private val LOG = logger()