Jetty的 WebSocket客户端只能接收到第一条消息
我在使用Jetty的 WebSocket服务器和WebSocket客户端时,遇到了一个问题:客户端只能收到来自服务器的第一条消息。
下面是我的主文件:
import jakarta.websocket.*
import jakarta.websocket.server.ServerEndpoint
import org.eclipse.jetty.ee10.servlet.ServletContextHandler
import org.eclipse.jetty.ee10.websocket.jakarta.server.config.JakartaWebSocketServletContainerInitializer
import org.eclipse.jetty.server.Server
class WebsocketServer {
private val server = Server(8080)
private val handler = ServletContextHandler("")
init {
server.handler = handler
JakartaWebSocketServletContainerInitializer.configure(handler) { _, container ->
container.addEndpoint(TestEndpoint::class.java)
}
}
fun start() {
server.start()
}
companion object {
@JvmStatic
fun main(args: Array<String>) {
WebsocketServer().start()
}
}
}
@ServerEndpoint("/websocket")
class TestEndpoint {
@OnOpen
fun setup() {
println("Websocket opened")
}
@OnMessage
fun respond(message: String, session: Session) {
session.basicRemote.sendText("Thanks for sending message: $message")
}
@OnError
fun onError(exception: Throwable, session: Session) {
println("Error: $exception")
}
@OnClose
fun onClose(session: Session) {
session.basicRemote.sendText("Goodbye")
}
}
下面是我的测试文件:
import org.eclipse.jetty.client.HttpClient
import org.eclipse.jetty.websocket.api.Callback
import org.eclipse.jetty.websocket.client.WebSocketClient
import org.junit.jupiter.api.Test
import java.net.URI
import java.util.concurrent.Executors
class WebsocketsTest {
@Test
fun test() {
Executors.newSingleThreadExecutor()
.use { executor ->
executor.submit {
WebsocketServer.main(emptyArray())
}
val websocketClient = WebSocketClient(HttpClient()).apply { start() }
val receivedMessages = mutableListOf<String>()
val listener = object : org.eclipse.jetty.websocket.api.Session.Listener {
override fun onWebSocketText(message: String) {
receivedMessages.add(message)
}
}
val session = websocketClient.connect(listener, URI.create("ws://localhost:8080/websocket"))
.get()
session.sendText("Hello World", Callback.NOOP)
session.sendText("Another message", Callback.NOOP)
Thread.sleep(1000)
session.close()
assert(receivedMessages.size == 2)
}
}
}
我正在使用以下库:
org.eclipse.jetty:jetty-server:12.1.10
org.eclipse.jetty.ee10:jetty-ee10-servlet:12.1.10
org.eclipse.jetty.ee10.websocket:jetty-ee10-websocket-jakarta-server:12.1.10
org.eclipse.jetty.websocket:jetty-websocket-jetty-client:12.1.10
如果你运行测试并调试它,你会发现服务器的第一条消息会被监听器接收到,但第二条消息却收不到。我已经深入研究了库的实现,试图找出消息丢失的原因,但仍然没能弄清楚。我希望我遇到的问题只是一个显而易见、可以轻易解决的问题。
解决方案
Jetty的 WebSocket实现要求你在接收任何新消息之前,显式地请求下一条消息。我不太清楚背后的原因,但我认为这是为了避免监听器代码因新消息可能随时异步到达而导致的潜在同步问题。
现在,默认情况下,监听器实现会在 open 事件之后自动请求第一条消息,因此你会收到那条消息。
然而,随后的消息需要通过显式调用会话的 demand() 方法来请求。请注意,这一步需要在 onWebSocketText() 中完成(即在消息被 received 之后),而不是在 session.sendText(...) 之后(即在消息被 sent 之后)。
所以你首先需要实现监听器的 onWebSocketOpen,在监听器内部保存会话以便后续使用,然后在那里也调用 demand()(因为你不再使用监听器的默认实现)。这看起来大概像这样:
val listener = object : org.eclipse.jetty.websocket.api.Session.Listener {
private var session: Session? = null
override fun onWebSocketOpen(session: Session?) {
this.session = session
session?.demand()
}
override fun onWebSocketText(message: String) {
receivedMessages.add(message)
session?.demand()
}
}
现在,每当接收到新消息时,你就会请求下一条,监听器的回调将再次被调用。现在你的测试通过了。
一种简单的替代方案是改用子接口 Session.Listener.AutoDemanding。这样连线就会自动完成,你的监听代码可以保持原状,如以前一样:
val listener = object : org.eclipse.jetty.websocket.api.Session.Listener.AutoDemanding {
override fun onWebSocketText(message: String) {
receivedMessages.add(message)
}
}
请注意,这需要你的 onWebSocketText 方法不要调用任何异步代码,例如向服务器再发送一条消息。更多信息,请参阅文档中的相关章节。
一个更好的替代方案是将WebSocket实现从基于Java的 Jetty转向用Kotlin原生编写的实现。想到的 Ktor框架 拥有丰富的特性集和更清晰的API界面。由于它原生就对协程友好,你不需要像Jetty那样费尽周折去把它做对。