refactor [lsp]: remove ktor dependency from lsp.protocol

The `lsp.protocol` module now doesn't depend on ktor.

All code that was depending on it was pushed in the module users.

GitOrigin-RevId: f065946f470ab9d13c457a0623267c2cb55dd0e5
This commit is contained in:
Ludwig Valda Vasquez
2025-11-20 16:02:36 +00:00
committed by intellij-monorepo-bot
parent 9040721730
commit 24604882cb
7 changed files with 1 additions and 207 deletions
-1
View File
@@ -27,7 +27,6 @@ jvm_library(
"@lib//:kotlinx-serialization-json",
"//fleet/util/core",
"@lib//:jetbrains-annotations",
"//fleet/ktor/network/tls",
"@lib//:kotlinx-io-core",
]
)
@@ -36,7 +36,6 @@
<orderEntry type="library" name="kotlinx-serialization-json" level="project" />
<orderEntry type="module" module-name="fleet.util.core" />
<orderEntry type="library" name="jetbrains-annotations" level="project" />
<orderEntry type="module" module-name="fleet.ktor.network.tls" />
<orderEntry type="library" name="kotlinx-io-core" level="project" />
</component>
</module>
@@ -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
}
@@ -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()
}
@@ -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"
@@ -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()
}
}
}
}
}
@@ -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<LspClient>()