iT邦幫忙

2026 iThome 鐵人賽

DAY 30
0
Software Development

Kotlin Ktor 實戰 101系列 第 30 篇

Kotlin Ktor 實戰 101 Day 30 WebSocket 即時推播

  • 分享至 

  • xImage
  •  

https://ithelp.ithome.com.tw/upload/images/20260909/20121948LfQo3aVCtG.jpg

這篇要來裝 ktor-server-websockets,把 Todo 的變更做成事件推出去,讓連著的人不用重新整理就看得到,順便看 day 26 到 day 28 那套認證授權放到一條長連線上還剩下多少效力,一條 socket 從 handshake 到關掉之間,它到底被檢查了幾次

這篇要完成什麼

  • 加上一行 ktor-server-websockets 相依
  • 事件用 sealed interface 加 @SerialName,序列化出來的 discriminator 是 type
  • webSocket { } 的 block 是一個活著的 session,collect 掛在那裡直到連線關掉
  • WebSocketOptions 那 4 個欄位翻原始碼看預設值,其中 2 個寫了等於沒寫
  • 推播對 HTTP 那一半的侵入只有 3 行加一個參數,查詢的 2 條路由一個字都沒動
  • handshake 就是一個普通的 HTTP GET,沒帶 token 的 401 跟打 /todos 拿到的一模一樣
  • 瀏覽器沒辦法在 handshake 上帶 header,token 只能走 query string
  • 連線活得比開它的那張 ticket 久,這篇沒有做主動關掉過期連線的機制
  • 2 個訊息一模一樣的 timeout,一個是必現的 receiver 遮蔽,一個是要跑 20 幾次才碰得到一次的 race

加上相依

build.gradle.kts 的 dependencies { } 加一行

     implementation(ktorLibs.server.swagger)
+    implementation(ktorLibs.server.websockets)

測試那邊要一個 client 的對應品

     testImplementation(ktorLibs.client.contentNegotiation)
+    testImplementation(ktorLibs.client.websockets)

事件長什麼形狀

先決定要送出去的東西,新增一個檔案 src/main/kotlin/com/cashwu/todo/TodoEvents.kt

package com.cashwu.todo

import kotlinx.coroutines.channels.BufferOverflow
import kotlinx.coroutines.flow.MutableSharedFlow
import kotlinx.coroutines.flow.SharedFlow
import kotlinx.coroutines.flow.asSharedFlow
import kotlinx.serialization.SerialName
import kotlinx.serialization.Serializable

@Serializable
sealed interface TodoEvent {

    @Serializable
    @SerialName("created")
    data class Created(val todo: Todo) : TodoEvent

    @Serializable
    @SerialName("updated")
    data class Updated(val todo: Todo) : TodoEvent

    @Serializable
    @SerialName("deleted")
    data class Deleted(val id: Int) : TodoEvent
}

class TodoEvents(bufferCapacity: Int = 64) {
    private val stream = MutableSharedFlow<TodoEvent>(
        replay = 0,
        extraBufferCapacity = bufferCapacity,
        onBufferOverflow = BufferOverflow.DROP_OLDEST,
    )

    val events: SharedFlow<TodoEvent> = stream.asSharedFlow()

    val subscriberCount: Int
        get() = stream.subscriptionCount.value

    suspend fun publish(event: TodoEvent) {
        stream.emit(event)
    }
}

sealed interface 加 @SerialName 走的是 kotlinx.serialization 的 closed polymorphism,3 個子類型都在同一個檔案裡,編譯期就知道總共有哪幾種,序列化出來的 discriminator 預設是 type,3 種事件長這樣

{"type":"created","todo":{"id":4,"title":"倒垃圾","done":false,"created_at":"2026-08-31T14:58:47.520890Z"}}
{"type":"updated","todo":{"id":1,"title":"買豆漿","done":true,"created_at":"2026-08-27T08:00:00Z"}}
{"type":"deleted","id":2}

created_at 那個欄位名還是 day 13 那個 @SerialName 的結果,Todo 這個型別在 HTTP 回應跟 WebSocket 事件裡是同一份,所以 client 解 todo 那一層用的可以是解 GET /todos 那份回應的同一段程式碼,欄位不會有 2 套名字

MutableSharedFlow 那 3 個參數各自是一個決定,replay = 0 表示新訂閱者不會拿到任何舊事件,它帶來一個真的 race,只是比想像中罕見,後面查那個 timeout 的時候會回來講,extraBufferCapacity 給 64、onBufferOverflow 是 DROP_OLDEST,也就是有訂閱者跟不上的時候丟掉最舊的那筆,而不是讓 emit 卡住等它,這個組合之下 emit 永遠不會真的 suspend,publish 掛著 suspend 是因為 emit 的簽章就是 suspend 的,不是因為它真的會等

subscriberCount 那個屬性不是功能,是為了測試才開的,理由在後面查 timeout 那一節

WebSocket 路由

新增一個檔案 src/main/kotlin/com/cashwu/todo/TodoSocket.kt

package com.cashwu.todo

import io.ktor.server.auth.principal
import io.ktor.server.routing.Route
import io.ktor.server.routing.openapi.hide
import io.ktor.server.websocket.WebSockets
import io.ktor.server.websocket.webSocket
import io.ktor.utils.io.ExperimentalKtorApi
import io.ktor.websocket.CloseReason
import io.ktor.websocket.Frame
import io.ktor.websocket.close
import kotlinx.coroutines.flow.collect
import kotlinx.serialization.json.Json

const val EVENTS_PATH = "/todos/events"

private val eventJson = Json { encodeDefaults = true }

@OptIn(ExperimentalKtorApi::class)
fun Route.todoEventsRoute(events: TodoEvents) {
    webSocket(EVENTS_PATH) {
        val user = call.principal<TodoUser>()
        if (user == null) {
            close(CloseReason(CloseReason.Codes.VIOLATED_POLICY, "沒有身分"))
            return@webSocket
        }
        events.events.collect { event ->
            send(Frame.Text(eventJson.encodeToString(TodoEvent.serializer(), event)))
        }
    }.hide()
}

fun webSocketOptions(): WebSockets.WebSocketOptions.() -> Unit = {
    pingPeriodMillis = 15_000
    timeoutMillis = 15_000
    maxFrameSize = 64L * 1024
    masking = false
}

webSocket { } 這個 block 的接收者是一個 session,裡面是 suspend 的,collect 會一直掛在那裡直到連線關掉,這就是 summary 裡「把連線放進協程模型」的意思,一條連線就是一個 coroutine,不用自己管 thread,也不用自己維護一份「現在有誰連著」的清單,那份清單就是 MutableSharedFlow 的訂閱者

.hide() 是 day 29 那個 @ExperimentalKtorApi 的函式,接一個 Route 回一個 Route,讓這條路由不要出現在 OpenAPI 規格裡,WebSocket 不是 OpenAPI 描述得了的東西,路徑掛在 paths 底下只會讓讀規格的人多困惑一次

close(CloseReason(...)) 那條分支在 authenticate 底下其實走不到,因為沒有身分的請求在 handshake 就被擋掉了,留著是因為 principal<T>() 的回傳型別是 nullable,day 27 那篇拆過什麼條件下會拿到 null,validate 回的型別跟這裡取的型別對不上就是其中一種

設定那邊,WebSockets 這個 plugin 的設定型別是 WebSockets.WebSocketOptions,下面這段不是要新增的檔案,是翻 ktor-server-websockets 自己的 WebSockets.kt 抄出來的,它是 WebSockets 裡面的一個巢狀類別,這一篇動到的 4 個欄位節錄出來是

public class WebSockets private constructor(
    // ...
) : CoroutineScope {

    @KtorDsl
    public class WebSocketOptions {

        public var pingPeriodMillis: Long = PINGER_DISABLED

        public var timeoutMillis: Long = 15_000L

        public var maxFrameSize: Long = Long.MAX_VALUE

        public var masking: Boolean = false

        // ...
    }
}

省略掉的是 contentConverter 這個欄位跟 channels { }、extensions { } 這 2 個函式,還有每個欄位上面那一大段 KDoc,這一篇都沒有動到

4 個都是 var,預設值就寫在宣告上,其中 PINGER_DISABLED 是 ktor-websockets 的一個常數,宣告是 public const val PINGER_DISABLED: Long = 0

所以 4 個裡面有 2 個寫了等於沒寫,timeoutMillis 本來就是 15000,那一行沒有改變任何東西,masking 本來就是 false,那一行也一樣,真的改到值的只有 pingPeriodMillis,從 0 也就是不送 ping 變成 15000,還有 maxFrameSize,從 Long.MAX_VALUE 也就是 9223372036854775807 降到 65536

這 4 個欄位這一輪都只是設了值,沒有測試過任何一個的行為,所以這篇不替那些數字的效果背書,先知道位置在哪裡,抄設定範例的時候順手把預設值再寫一次是很常見的事,翻一下原始碼花的時間比爭論便宜

接線

src/main/kotlin/com/cashwu/todo/Application.kt 的 DI 容器多註冊一個,底下取用的地方也多一行

    dependencies {

        // ...

+       provide<TodoEvents> { TodoEvents() }
    }

    // ...

+   val events: TodoEvents by dependencies

同一個檔案裝 plugin 的地方

+   install(WebSockets) {
+       webSocketOptions()()
+   }

    // ...

webSocketOptions()() 那 2 對括號不是筆誤,前面那對是呼叫那個回傳 lambda 的函式,後面那對是把拿到的 lambda 套用在 WebSocketOptions 上

routing { } 裡面,新的路由跟原本的 CRUD 放在同一個 authenticate 底下

       authenticate(TODO_AUTH) {
           meRoute()
+          todoEventsRoute(events)
-          todoRoutes(repository)
+          todoRoutes(repository, events)
       }

src/main/kotlin/com/cashwu/todo/TodoRoutes.kt 的簽章多一個參數

-fun Route.todoRoutes(repository: TodoRepository) {
+fun Route.todoRoutes(repository: TodoRepository, events: TodoEvents) {

同一個檔案裡 3 個會改變資料的 handler 各多一行,post { } 這條

       post {
           val request = call.receive<CreateTodoRequest>()
           val todo = repository.create(request.title, request.done)
+          events.publish(TodoEvent.Created(todo))
           call.respond(HttpStatusCode.Created, todo)
       }.describeJson(

put("{id}") { } 這條

       put("{id}") {
           val id = call.todoId()
           val request = call.receive<UpdateTodoRequest>()
           val todo = repository.update(id, request.title, request.done)
               ?: throw ApiException(HttpStatusCode.NotFound, "找不到 id $id 的待辦")
+          events.publish(TodoEvent.Updated(todo))
           call.respond(todo)
       }.describeJson(

delete { } 這條,它在 day 28 那個 authorize(ADMIN_ROLE) 裡面

       delete {
           val id = call.todoId()
           if (!repository.delete(id)) {
              throw ApiException(HttpStatusCode.NotFound, "找不到 id $id 的待辦")
           }
+          events.publish(TodoEvent.Deleted(id))
           call.respond(HttpStatusCode.NoContent)
       }.describeJson(

3 行都排在 respond 前面,寫成這樣的理由是順序反過來的話,事件送不出去的時候 client 已經拿到成功的回應了,那時候沒得補救

整個推播機制對 HTTP 這一半的侵入就是這 3 行加一個參數,查詢的 2 條路由一個字都沒動,因為它們不改資料

handshake 就是一個普通的 HTTP GET

這一段讓 day 26 到 day 28 那 3 篇原封不動地生效

WebSocket 連線的第 1 步是一個 HTTP GET,帶著 Upgrade 跟 Sec-WebSocket-Key 那幾個 header,server 認可了才回 101 換協定,既然是 GET,Ktor 的認證就照樣攔得到

這一節到後面幾節的請求都要帶 token,所以先登入一次把 access token 存進 shell 變數,做法跟 day 27、day 28 一樣

curl -s -X POST localhost:8080/login \
  -H 'Content-Type: application/json' \
  -d '{"name":"alice","password":"alice-secret"}' > /tmp/login.json
AT=$(python3 -c 'import json; print(json.load(open("/tmp/login.json"))["accessToken"])')

然後用 curl 手動組一個 handshake,帶著 Authorization

curl -i -s --http1.1 -H "Connection: Upgrade" -H "Upgrade: websocket" \
  -H "Sec-WebSocket-Version: 13" -H "Sec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==" \
  -H "Authorization: Bearer $AT" localhost:8080/todos/events --max-time 3

回來的是

HTTP/1.1 101 Switching Protocols
X-Request-Id: ow5/wbzhoiy6
X-Response-Time: 20ms
Upgrade: websocket
Connection: Upgrade
Sec-WebSocket-Accept: s3pPLMBiTxaQ9kYGzzhZRbK+xOo=

同一個指令把 Authorization 拿掉

HTTP/1.1 401 Unauthorized
X-Request-Id: c9p8lv-/-7jo
X-Response-Time: 3ms
WWW-Authenticate: Bearer realm="todo-api"
Content-Length: 71
Content-Type: application/json

{"status":401,"message":"請帶著有效的 token 再來","details":[]}

這個 401 上面每一行都是前面幾篇留下來的,WWW-Authenticate 是 day 27 的 challenge、body 是 day 28 的 respondFixedJson、X-Request-Id 是 day 16、X-Response-Time 是 day 10,外面那層 authenticate(TODO_AUTH) 跟 TodoUser 是 day 26 的,WebSocket 這一層什麼都沒有多做

101 那個回應也一樣掛著那 2 個 header,因為 handshake 走的就是它們本來在處理的那條 HTTP 路徑,差別在 101 之後,那條連線不再經過 pipeline

瀏覽器帶不了 header

上面那個做法有一個前提,你的 client 要能自己組 header,new WebSocket(url) 只收 URL 跟 subprotocol,沒有地方放 header,所以瀏覽器用不了 Authorization 這條路

Ktor 的 jwt provider 有 authHeader 這個 hook,day 26 那篇翻 bearer 的 Config 時就看到過它,當時的結論是「可以換掉從哪裡讀 token 的規則」

這篇就使用到了,src/main/kotlin/com/cashwu/todo/Auth.kt 先多一個常數

 const val ROLES_CLAIM = "roles"
+const val ACCESS_TOKEN_PARAM = "access_token"

同一個檔案的 jwt(TODO_AUTH) 區塊裡加上 hook,第 1 版是這樣寫的

     jwt(TODO_AUTH) {
         realm = config.realm
+        authHeader { call ->
+            val header = call.request.authorization() ?: call.request.queryParameters[ACCESS_TOKEN_PARAM]
+                ?.let { "Bearer $it" }
+            header?.let { parseAuthorizationHeader(it) }
+        }

         // ...
      }

第 1 版的邏輯是 header 優先,沒有 header 才去看 query string,撿到的話自己拼成 Bearer xxx 再交給 Ktor 解析,這裡還少了一條 RFC 6750 2 的限制,client 不能在同一個請求同時使用 2 種傳遞 token 的方式,最後的版本因此不再默默挑一個,而是讓同時帶 header 與 query token 的請求驗證失敗,實測只用 query string 打 handshake

curl -i -s --http1.1 -H "Connection: Upgrade" -H "Upgrade: websocket" \
  -H "Sec-WebSocket-Version: 13" -H "Sec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==" \
  "localhost:8080/todos/events?access_token=$AT" --max-time 3
HTTP/1.1 101 Switching Protocols
X-Request-Id: llxjzx0tfa7o
Cache-Control: private
X-Response-Time: 3ms
Upgrade: websocket
Connection: Upgrade
Sec-WebSocket-Accept: s3pPLMBiTxaQ9kYGzzhZRbK+xOo=

那行 Cache-Control: private 是最後的版本才加上去的,理由在下面 2.3 那 2 條 SHOULD 一起講

curl 有自己組 header 的自由,瀏覽器沒有,所以這一節開頭那句話這一輪拿真的瀏覽器驗了一次

下面這 2 行是另外開一個簡易的前端頁面跑的,不是 todo-api 的一部分,at 就是上面那張 access token

const ws = new WebSocket(`ws://localhost:8080/todos/events?access_token=${at}`)
ws.onmessage = (e) => console.log(e.data)

連上之後在終端機分別打 POST /todos、PUT /todos/1、DELETE /todos/2,console 上依序印出來的是

{"type":"created","todo":{"id":4,"title":"倒垃圾","done":false,"created_at":"2026-09-04T22:32:26.235079Z"}}
{"type":"updated","todo":{"id":1,"title":"買豆漿","done":true,"created_at":"2026-08-27T08:00:00Z"}}
{"type":"deleted","id":2}

這 3 個 frame 跟前面第 3 節列的形狀是同一份,後 2 個更是跟 TodoSocketTest 那 2 個整串比對的斷言逐字相同,差別只在這次是瀏覽器收的

access_token 這個參數名不是隨便取的,RFC 6750 2.3 定義了它

2.3.  URI Query Parameter

   When sending the access token in the HTTP request URI, the client
   adds the access token to the request URI query component as defined
   by "Uniform Resource Identifier (URI): Generic Syntax" [RFC3986],
   using the "access_token" parameter.

同一節接下來那段是勸退

   Because of the security weaknesses associated with the URI method
   (see Section 5), including the high likelihood that the URL
   containing the access token will be logged, it SHOULD NOT be used
   unless it is impossible to transport the access token in the
   "Authorization" request header field or the HTTP request entity-body.
   Resource servers MAY support this method.

2 件事要一起看

那個 SHOULD NOT 有例外條款,unless it is impossible to transport the access token in the "Authorization" request header field or the HTTP request entity-body,那是 2 個選言,兩邊都要關掉例外才成立,瀏覽器的 WebSocket API 沒有地方放 header 是一邊,另一邊是 handshake 本來就是 GET,沒有 entity-body 可以放,兩邊都關掉了,所以這個 API 落在例外裡面,不是「懶得放 header」

它給的理由是 the high likelihood that the URL containing the access token will be logged,而這個 API 的 log 不會,day 16 那篇之前就處理過同一件事,當時發現 day 15 補的那行 log.error 用 uri 把 ?token=super-secret 原封不動寫進了 log,改成 call.request.path() 之後 query string 不再進來,而 CallLogging 本來用的就是 path(),那篇留下的測試是

fun `no log line carries the query string`() = testApplication {
    probeApp {
        get("/boom") { throw RuntimeException("炸了") }
    }

    client.get("/boom?token=super-secret")
    awaitRequestLines(1)

    assertTrue(events().none { it.formattedMessage.contains("super-secret") })
    }

那個測試驗證的是 log 這一層的行為,跟這篇的決定隔了 14 篇,當初也不是為了 WebSocket 寫的,這一輪在真的 server 上再驗一次,server 是用 DB_PASSWORD=todo ./gradlew run > /tmp/server.log 2>&1 起的,所以整份 log 都在檔案裡

curl -s -o /dev/null "localhost:8080/todos?access_token=$AT"
grep -c "$AT" /tmp/server.log

那個 token 在整份 log 裡是

0

而那一行 access log 是

22:48:42.364 INFO  [mp6rlfbfu8oj] Application -- 200 GET /todos 2ms

只有路徑,沒有 query string

2.3 在勸退那段之前還有 2 條 SHOULD,一條給 client 一條給 server

   Clients using the URI Query Parameter method SHOULD also send a
   Cache-Control header containing the "no-store" option.  Server
   success (2XX status) responses to these requests SHOULD contain a
   Cache-Control header with the "private" option.

server 那一半這篇做了,理由跟 log 那件事同一條線,token 在 URL 裡的時候,中間任何一個會存東西的環節都是外洩的位置,快取也算一個,authHeader 因此改寫成 2 條路,走 query string 的那條順手加一個 response header,變成下面這樣

    authHeader { call ->
        val fromHeader = call.request.authorization()?.let(::parseAuthorizationHeader)
        val fromQuery = call.request.queryParameters[ACCESS_TOKEN_PARAM]
        if (fromHeader != null && fromQuery != null) return@authHeader null
        if (fromHeader != null) return@authHeader fromHeader
        if (fromQuery == null) return@authHeader null
        call.response.header(HttpHeaders.CacheControl, "private")
        parseAuthorizationHeader("Bearer $fromQuery")
    }

拆成幾條 early return 之後,「哪一條路進來的」在程式碼上就分得開了,第 1 版那個 ?: 串起來的寫法沒有位置處理重複憑證,也沒有地方只替 query token 加 response header,實測帶 query token 的回應,2 塊輸出都只留 header 的部分

curl -i -s "localhost:8080/todos?access_token=$AT"
HTTP/1.1 200 OK
X-Request-Id: if8/3y8intrr
Cache-Control: private
X-Response-Time: 5ms
Content-Length: 245
Content-Type: application/json

帶 header 的那次沒有那一行

curl -i -s -H "Authorization: Bearer $AT" localhost:8080/todos
HTTP/1.1 200 OK
X-Request-Id: nlxlycoi4h1k
X-Response-Time: 1ms
Content-Length: 245
Content-Type: application/json

2 種進來的方式回應不一樣,而 client 那一半的 Cache-Control: no-store 是 client 的責任,這篇管不到

day 27 的 challenge 也要同步看這 2 個來源,變成下面這樣

challenge { defaultScheme, realm ->
    val parsedHeader = call.request.authorization()?.let(::parseAuthorizationHeader)
    val bearerHeader = parsedHeader?.authScheme?.equals(defaultScheme, ignoreCase = true) == true
    val queryToken = call.request.queryParameters[ACCESS_TOKEN_PARAM]
    val parameters = when {
        parsedHeader != null && queryToken != null ->
            mapOf(HttpAuthHeader.Parameters.Realm to realm, "error" to "invalid_request")

        bearerHeader || queryToken != null ->
            mapOf(HttpAuthHeader.Parameters.Realm to realm, "error" to "invalid_token")

        else -> mapOf(HttpAuthHeader.Parameters.Realm to realm)
    }
    call.respond(
        UnauthorizedResponse(
            HttpAuthHeader.Parameterized(defaultScheme, parameters, HeaderValueEncoding.QUOTED_ALWAYS)
        )
    )
}

只帶 Bearer header 或只帶 query token、但內容驗證失敗時回 error="invalid_token",兩邊都帶則回 error="invalid_request",如果兩邊都沒有,或只送來不支援的認證方式,就保留 realm,這樣才不會把單獨的 Basic ... 誤報成壞掉的 Bearer token,也不會讓重複憑證靠優先順序偷偷通過,兩邊都帶的那次實測是

curl -i -s "localhost:8080/todos?access_token=$AT" -H "Authorization: Bearer $AT"
HTTP/1.1 401 Unauthorized
X-Request-Id: wc5v+4p2aty0
X-Response-Time: 0ms
WWW-Authenticate: Bearer realm="todo-api", error="invalid_request"
Content-Length: 71
Content-Type: application/json

Cache-Control 那一行不在,因為 authHeader 認出 2 種憑證同時出現就直接回 null 了,沒有走到加 header 那幾行

有一個誠實的但書要補,RFC 那句限定的是 Server success (2XX status) responses,而 authHeader 是在還不知道結果的時候就跑了,所以只要 URL 帶了 ?access_token=,非 2XX 的回應也會拿到那個 header,實測一個壞 token 的 401

curl -i -s "localhost:8080/todos?access_token=not-a-token"
HTTP/1.1 401 Unauthorized
X-Request-Id: +tzt2xdpfzav
Cache-Control: private
X-Response-Time: 1ms
WWW-Authenticate: Bearer realm="todo-api", error="invalid_token"
Content-Length: 71
Content-Type: application/json

跟一個好 token 打不存在的 id 的 404

curl -i -s "localhost:8080/todos/999?access_token=$AT"
HTTP/1.1 404 Not Found
X-Request-Id: gygiqiit/dbx
Cache-Control: private
X-Response-Time: 5ms
Content-Length: 66
Content-Type: application/json

比 SHOULD 要求的寬,而且不會有害,帶著 token 的 URL 回的 401 跟 404 本來也不該進共享快取,但這件事說明了 authHeader 是「從哪裡讀憑證」的鉤子,拿它寫 response header 是超出語意的用法,副作用落在哪些回應上不是那個鉤子管得到的

做完這些之後 RFC 還是沒有給好臉色,同一節在勸退那句之後有一句比 SHOULD NOT 更重的

   This method is included to document current use; its use is not
   recommended, due to its security deficiencies (see Section 5) and
   also because it uses a reserved query parameter name, which is
   counter to URI namespace best practices, per "Architecture of the
   World Wide Web, Volume One" [W3C.REC-webarch-20041215].

its use is not recommended 這句沒有例外條款,所以這篇的立場不是「RFC 說可以」,是「RFC 說不建議,我們知道它為什麼不建議,其中一項理由在這個專案上不成立,該做的那條 SHOULD 做了,剩下的自己認」

但書有 2 個,第 1,authHeader 是掛在整個 provider 上的,所以 ?access_token= 對所有受保護的路由都有效,不只 WebSocket,這件事後面有一個測試把它寫成現況,第 2,log 不記不代表安全,反向代理、瀏覽器歷史、Referer 都還是會看到那個 URL,RFC 講的那些 weaknesses 只被拆掉了一項

這篇加的測試

todoRoutes() 的簽章雖然多了一個參數,但整個 src/test/ 底下沒有人直接呼叫它,day 19 換掉 repository 那次也遇過同一件事,測試要嘛走 HTTP、要嘛自己拼路由,這件事本身是一個結果,推播是加上去的,沒有改掉任何一條 HTTP 路由對外的行為

新增一個檔案 src/test/kotlin/com/cashwu/todo/TodoSocketTest.kt,底下每個測試都會用到 3 個 helper

class TodoSocketTest {

    private lateinit var events: TodoEvents

    private fun ApplicationTestBuilder.eventApplication() {
        todoApplication()
        application { events = dependencies.resolve<TodoEvents>() }
    }

    private suspend fun awaitSubscribers(count: Int) {
        withTimeout(5_000) {
            while (events.subscriberCount < count) {
                delay(10)
            }
        }
    }

    private suspend fun ApplicationTestBuilder.createTodo(title: String): HttpResponse =
        bearerClient().post("/todos") {
            contentType(ContentType.Application.Json)
            setBody("""{"title":"$title","done":false}""")
        }
}

eventApplication() 除了起 app,還把 DI 容器裡那個 TodoEvents 接到欄位上,測試看到的跟路由用的才是同一個實體,todoApplication() 是 TestApp.kt 裡那個把整個 app 起起來、回一個已經帶著 TEST_TOKEN 的 client 的 helper,day 29 也用同一個

awaitSubscribers(n) 等訂閱數到位再往下走,為什麼需要它是後面 replay = 0 那一節的事,它讀的 subscriberCount 就是 TodoEvents 開給測試的那個屬性

createTodo(title) 就是用 bearerClient() 打一次 POST /todos,bearerClient(token) 是 TestApp.kt 的,預設值是 TEST_TOKEN

第 1 個驗證的是最基本的那條路

@Test
fun `a subscriber gets the event a post produces`() = testApplication {
    eventApplication()
    val ws = createClient { install(WebSockets) }

    ws.webSocket(EVENTS_PATH, request = { header(HttpHeaders.Authorization, "Bearer $TEST_TOKEN") }) {
        awaitSubscribers(1)
        val created = createTodo("倒垃圾")
        assertEquals(HttpStatusCode.Created, created.status)

        val frame = withTimeout(5_000) { incoming.receive() } as Frame.Text
        val text = frame.readText()

        assertTrue(text.startsWith("""{"type":"created""""), text)
        assertTrue(text.contains(""""title":"倒垃圾""""), text)
    }
}

request = { header(...) } 那個參數就是前面 curl 手動組的那個 Authorization,Ktor 的 client 沒有瀏覽器那個限制,這條路走得通,斷言盯的是 2 件事,discriminator 是 created,todo 那一層帶著標題

第 2 個是更新跟刪除,這個測試踩過一次名實不符,第 1 版的名字叫 update and delete produce their own events,body 卻只做了一個 create,斷言也就只斷到 created 那個 frame,名字說的事情一件都沒測到,改成真的做那 2 件事之後長這樣

@Test
fun `update and delete produce their own events`() = testApplication {
    eventApplication()
    val ws = createClient { install(WebSockets) }

    ws.webSocket(EVENTS_PATH, request = { header(HttpHeaders.Authorization, "Bearer $TEST_TOKEN") }) {
        awaitSubscribers(1)
        val client = bearerClient()

        client.put("/todos/1") {
            contentType(ContentType.Application.Json)
            setBody("""{"title":"買豆漿","done":true}""")
        }
        val updated = (withTimeout(5_000) { incoming.receive() } as Frame.Text).readText()

        client.delete("/todos/2")
        val deleted = (withTimeout(5_000) { incoming.receive() } as Frame.Text).readText()

        assertEquals(
            """{"type":"updated","todo":{"id":1,"title":"買豆漿","done":true,"created_at":"2026-08-27T08:00:00Z"}}""",
            updated,
        )
        assertEquals("""{"type":"deleted","id":2}""", deleted)
    }
}

動的是測試資料裡的第 1 筆跟第 2 筆,2 個事件在同一條 socket 上依序收,這 2 個斷言跟上面那個不一樣,是整串比對不是 contains,所以前面第 3 節那 2 個 frame 的形狀全部被這個測試驗到了,discriminator 是 updated 跟 deleted、todo 那一層的欄位連順序都算數、Deleted 只帶一個 id 而不是整筆資料,哪天有人替 TodoEvent 換一個 discriminator 名字或是替 Todo 加一個欄位,這裡會失敗,而那正是 client 端會跟著壞掉的地方

第 3 個驗證的是沒帶 token 連不上

@Test
fun `the handshake without a token is rejected`() = testApplication {
    eventApplication()
    val ws = createClient { install(WebSockets) }

    var opened = false
    runCatching {
        ws.webSocket(EVENTS_PATH) { opened = true }
    }

    assertEquals(false, opened)
}

這裡用 runCatching 加一個旗標而不是斷言 status code,因為 client 那一層在 handshake 失敗時是丟例外的,拿不到 401 這個值,那個例外印出來是

IllegalStateException: WebSocket connection failed

訊息裡沒有 401,這也是為什麼前面那個 401 的證據是 curl 打出來的而不是測試斷出來的,這個測試能保證的是「連不上」,「回的是什麼」在 curl 那一節

第 4 個驗證的是 .hide()

@Test
fun `the events route stays out of the openapi spec`() = testApplication {
    todoApplication()

    val yaml = client.get("/$SPEC_PATH").bodyAsText()

    assertTrue(!yaml.contains("/todos/events"), yaml)
}

SPEC_PATH 是 day 29 那個常數,這個測試接的是 day 29 那份 OpenApiTest 的完整路徑清單斷言,那邊比對的是「規格裡有的路徑」,這邊驗證的是「這條不該進去」

第 5 個驗證的是 RFC 那條 Cache-Control 的 SHOULD

@Test
fun `a query string token gets the cache control the rfc asks for`() = testApplication {
    todoApplication()

    val fromQuery = client.get("/todos?$ACCESS_TOKEN_PARAM=$TEST_TOKEN")
    val fromHeader = bearerClient().get("/todos")

    assertEquals("private", fromQuery.headers[HttpHeaders.CacheControl])
    assertEquals(null, fromHeader.headers[HttpHeaders.CacheControl])
}

2 個斷言是一對,第 1 個確認走 query string 的回應有那個 header,第 2 個確認走 Authorization 的沒有,因為那條路沒有這個需求,多送一個 header 只是雜訊,少了第 2 個斷言的話,把 header 那一行搬到 authHeader 最外層也會通過,而那是錯的實作

第 6 個把前面那個但書寫成現況

@Test
fun `query tokens work on ordinary routes but cannot duplicate the header`() = testApplication {
    todoApplication()

    val response = client.get("/todos?$ACCESS_TOKEN_PARAM=$TEST_TOKEN")
    val duplicate = bearerClient().get("/todos?$ACCESS_TOKEN_PARAM=$TEST_TOKEN")

    assertEquals(HttpStatusCode.OK, response.status)
    assertEquals(HttpStatusCode.Unauthorized, duplicate.status)
    assertContains(duplicate.headers[HttpHeaders.WWWAuthenticate].orEmpty(), "error=\"invalid_request\"")
}

authHeader 掛在 provider 上,所以 query token 能打開普通路由是目前的行為,同一個測試也驗了「不能再帶一份 Authorization header」,哪天有人決定只讓 WebSocket 那條路吃 query string,第 1 個斷言會失敗,那個失敗是提醒不是壞消息

剩下的 4 個裡有 2 個一句話就講完,query string 開得起 WebSocket、沒有訂閱者的時候 publish 就是沒有人接,另外 2 個各自帶出一件要多講的事

連線活得比開它的那張 ticket 久

認證只在 handshake 發生一次,之後那條 socket 就不再檢查任何東西

測試是這樣寫的,bearerClient(shortLived) 用的是那張短命的 ticket,不是預設的 TEST_TOKEN

@Test
fun `the connection outlives the token that opened it`() = testApplication {
    eventApplication()
    val shortLived = TokenIssuer(todoTestConfig.auth.jwt.copy(accessTtlSeconds = 1), Clock.systemUTC())
        .accessToken(TEST_USER)
    val ws = createClient { install(WebSockets) }

    ws.webSocket("$EVENTS_PATH?$ACCESS_TOKEN_PARAM=$shortLived") {
        awaitSubscribers(1)
        delay(2_000)

        assertEquals(
            HttpStatusCode.Unauthorized,
            bearerClient(shortLived).get("/todos").status,
        )

        createTodo(" ticket 過期了連線還在")
        val frame = (withTimeout(5_000) { incoming.receive() } as Frame.Text).readText()

        assertTrue(frame.contains(" ticket 過期了連線還在"), frame)
    }
}

拿到一張 TTL 一秒的 token 開連線,等 2 秒,然後在同一個測試裡做 2 件事,同一張 ticket 打 /todos 拿到 401,這是 day 27 那個 verifier 在做事,同一張 ticket 開的那條 socket 照樣收到了新事件,這個測試通過

這不是 Ktor 的 bug,是 WebSocket 的模型本來就這樣,HTTP 的每一個請求都會重新經過認證,而 socket 只有 handshake 那一次會,要處理的話得自己做,例如在 session 裡記下 exp、用 withTimeout 或是定期比對時間主動 close,這篇沒有做,留成債

兩個長得一模一樣的 timeout

TodoSocketTest 第 1 版,2 個訂閱者那個測試 timeout 了

kotlinx.coroutines.TimeoutCancellationException: Timed out waiting for 5000 ms

直覺是 replay = 0 的 race,畢竟那是這個設計裡最明顯的一個時序缺口,那個猜測是錯的

錯在哪裡是用 2 組對照實驗量出來的,把 val outer = this 拿掉、2 次都寫 incoming.receive(),跑 3 次,3 次都失敗,而且失敗的都是 two subscribers both see the same event(),把 5 處 awaitSubscribers(n) 全部拿掉,跑 20 次,只有 1 次失敗,失敗的是 the connection outlives the token that opened it(),1 個必現,1 個要跑 20 幾次才碰得到 1 次,第 1 版那個 timeout 的成因顯然是前者

前者是 receiver 遮蔽,2 層巢狀的 webSocket { } 裡面,this@webSocket 指到的是最內層那個,所以外層的 session 要自己接成變數

@Test
fun `two subscribers both see the same event`() = testApplication {
    eventApplication()
    val first = createClient { install(WebSockets) }
    val second = createClient { install(WebSockets) }

    first.webSocket("$EVENTS_PATH?$ACCESS_TOKEN_PARAM=$TEST_TOKEN") {
        val outer = this
        second.webSocket("$EVENTS_PATH?$ACCESS_TOKEN_PARAM=$TEST_TOKEN") {
            awaitSubscribers(2)
            createTodo("2 個人都要收到")

            val inner = (withTimeout(5_000) { incoming.receive() } as Frame.Text).readText()
            val outerFrame = (withTimeout(5_000) { outer.incoming.receive() } as Frame.Text).readText()

            assertEquals(inner, outerFrame)
            assertTrue(inner.contains("2 個人都要收到"), inner)
        }
    }
}

沒有那個 outer,2 次 receive() 讀的是同一條 socket,第 1 次把唯一那個 frame 拿走,第 2 次就永遠等下去,這是 Kotlin 帶接收者的 lambda 疊起來的通例,只是在這裡的失敗長得像「事件沒送到第 2 個人」

而 replay = 0 的 race 也是真的,只是罕見

webSocket { } 的 block 開始執行跟 events.events.collect 真的登記成訂閱者之間有一小段空檔,MutableSharedFlow 的 replay 是 0,那段時間 emit 出去的事件對還沒訂閱的人就是不存在,20 次裡的那 1 次就是它

擋它的方法是等到訂閱數到位再動作,前面那個 awaitSubscribers(n) 就是為了它寫的,TodoEvents.subscriberCount 那個屬性也是

要分清楚的是這一步擋掉的只有測試的不穩定,正式環境那個 race 沒有被解掉,只是被看見了,一個 client 連上來之後馬上改一筆資料,有機會收不到自己那一筆,要解的話 replay 給一個大於 0 的值,代價是新連線會先收到舊事件,那對「即時推播」這個用途不一定是想要的行為

這一段真正的教訓是查錯的方向,2 個問題丟出來的訊息一模一樣,都是 Timed out waiting for 5000 ms,一個必現一個罕見,而第 1 次查的時候挑了比較有趣的那個解釋,加了 awaitSubscribers 之後測試還是失敗,才回頭找到必現的那一個,day 27 也有過一次類似的,那次是測試通過但通過的理由不對,擋下 token 的是 iat 的檢查不是過期,方向剛好相反,同樣是「訊息對上了就以為原因對了」

跟 Relix 的對照

Relix 沒有 WebSocket,day 31 那份「刻意不做的東西」清單第 1 條就是它,所以這篇沒有對照組,只有一句話可以接著講

- WebSocket,需要完全不同的 connection 管理模型

那句話這篇可以具體化,RelixHandler 是 suspend RelixCall.() -> RelixResponse,收一個 call 回一個 response,一次呼叫就結束,而 webSocket { } 是一個活著的 session,collect 掛在那裡不回來,簽章上根本沒有「回一個 response」的位置,不是加一個 handler 型別就塞得進去的

而且那個「刻意」有點客氣,Relix 用的是 JDK HttpServer,那個 API 本身就沒有升級到 WebSocket 的路,在那個 engine 上是做不到,engine 選什麼會決定框架的天花板在哪裡


小結

這篇的推播骨架有 3 個核心,一個 MutableSharedFlow 包成 TodoEvents、一條 webSocket 路由、3 個改資料的 handler 各多一行 publish,ktor-server-websockets

事件用 sealed interface 加 @SerialName 走 closed polymorphism,discriminator 預設是 type,Todo 在 HTTP 回應跟事件裡是同一份,client 解事件那一層用的可以是解 GET /todos 的同一段程式碼

webSocket { } 的 block 是一個活著的 session,collect 掛在裡面直到連線關掉,一條連線就是一個 coroutine,訂閱者清單就是那個 SharedFlow 自己,路由後面接 day 29 那個 .hide() 不讓它進 OpenAPI 規格,而 WebSocketOptions 動到的 4 個欄位有 2 個寫了等於沒寫,timeoutMillis 本來就是 15000、masking 本來就是 false

推播對 HTTP 那一半的侵入是 3 行 publish 加一個參數,排在 respond 前面是因為順序反過來的話,事件送不出去時 client 已經拿到成功的回應了

handshake 就是一個普通的 HTTP GET,所以 day 26 到 day 28 那套認證授權原封不動生效,帶著 Authorization 回的是 101,不帶回的 401 跟打 /todos 拿到的一模一樣,WebSocket 這一層什麼都沒有多做

瀏覽器的 new WebSocket(url) 沒有地方放 header,token 只能走 query string,做法是 authHeader 鉤子讀 RFC 6750 §2.3 定義的 access_token,兩邊都帶則以 invalid_request 拒絕,同一節說它 SHOULD NOT 但有例外條款,瀏覽器沒得放 header、handshake 又是沒有 entity-body 的 GET,2 個選言都關掉了,該做的那條 SHOULD 也做了,走 query string 進來的回應帶 Cache-Control: private,即使這樣 RFC 還有一句沒有例外的 its use is not recommended,所以立場是知道它為什麼不建議,不是「RFC 說可以」

第 1 版那個 Timed out waiting for 5000 ms 查錯了方向,必現的成因是 receiver 遮蔽,2 層巢狀的 webSocket { } 裡 this@webSocket 指到最內層,外層的 session 要自己接成 outer,而 replay = 0 的 race 也是真的,只是 20 次才碰到 1 次,2 個問題的訊息一模一樣,挑了比較有趣的那個解釋就查錯邊

這篇留下的問題有 5 個,連線活得比 token 久,沒有主動關掉過期連線的機制,那是這幾筆裡最該補的一個,?access_token= 對所有受保護的路由都有效而不只 WebSocket,要縮小範圍的話 authHeader 那一層要自己看路徑,replay = 0 的 race 沒有解,剛連上就改資料有機會收不到自己那一筆,事件只在單一個 process 裡廣播,多台機器的話 MutableSharedFlow 要換成外部的 pub/sub,那是整個機制唯一沒辦法靠這個檔案解決的一筆,最後是 WebSocketOptions 那 4 個欄位,這一輪都只是設了值,沒有量過任何一個的行為,其中 2 個還是把預設值再寫了一次


下一篇

下一篇又換到另一邊,到這裡 todo-api 一直是被打的那一方,day 31 做 HttpClient 跟 MockEngine,讓這個服務去打別人,順便處理「測試怎麼不依賴外部服務」這件事,這篇的 client 端其實已經先用過一次了,createClient { install(WebSockets) } 那個 ktor-client-websockets 就是 Ktor client 的一部分


參考資料


同步刊登於 Blog

圖片來源:AI 產生


上一篇
Kotlin Ktor 實戰 101 Day 29 OpenAPI 文件與 Swagger UI
系列文
Kotlin Ktor 實戰 101 共 30 篇
圖片
  熱門推薦
圖片
{{ item.channelVendor }} | {{ item.webinarstarted }} |
{{ formatDate(item.duration) }}
直播中

尚未有邦友留言

立即登入留言