-
Notifications
You must be signed in to change notification settings - Fork 1k
/
TestApplicationResponse.kt
160 lines (129 loc) · 4.96 KB
/
TestApplicationResponse.kt
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
/*
* Copyright 2014-2019 JetBrains s.r.o and contributors. Use of this source code is governed by the Apache 2.0 license.
*/
package io.ktor.server.testing
import io.ktor.http.*
import io.ktor.http.content.*
import io.ktor.server.engine.*
import io.ktor.server.request.*
import io.ktor.server.response.*
import io.ktor.util.*
import io.ktor.utils.io.*
import io.ktor.utils.io.charsets.*
import io.ktor.utils.io.core.*
import kotlinx.atomicfu.*
import kotlinx.coroutines.*
import kotlin.coroutines.*
/**
* A test call response received from a server.
* @property readResponse if response channel need to be consumed into byteContent
*/
public class TestApplicationResponse(
call: TestApplicationCall,
private val readResponse: Boolean = false
) : BaseApplicationResponse(call), CoroutineScope by call {
private val scope: CoroutineScope get() = this
private val timeoutAttributes get() = call.attributes.getOrNull(timeoutAttributesKey)
/**
* Gets a response body text content. Could be blocking. Remains `null` until response appears.
*/
public val content: String?
get() {
val charset = headers[HttpHeaders.ContentType]?.let { ContentType.parse(it).charset() } ?: Charsets.UTF_8
return byteContent?.let { charset.newDecoder().decode(ByteReadPacket(it)) }
}
/**
* Response body byte content. Could be blocking. Remains `null` until response appears.
*/
private val _byteContent = atomic<ByteArray?>(null)
public var byteContent: ByteArray?
get() = when {
_byteContent.value != null -> _byteContent.value
responseChannel == null -> null
else -> runBlocking { responseChannel!!.toByteArray() }
}
private set(value) {
_byteContent.value = value
}
private var responseChannel: ByteReadChannel? = null
private var responseJob: Job? = null
/**
* Get completed when a response channel is assigned.
* A response could have no channel assigned in some particular failure cases so the deferred could
* remain incomplete or become completed exceptionally.
*/
internal val responseChannelDeferred: CompletableJob = Job()
override fun setStatus(statusCode: HttpStatusCode) {}
override val headers: ResponseHeaders = object : ResponseHeaders() {
private val builder = HeadersBuilder()
override fun engineAppendHeader(name: String, value: String) {
builder.append(name, value)
}
override fun getEngineHeaderNames(): List<String> = builder.names().toList()
override fun getEngineHeaderValues(name: String): List<String> = builder.getAll(name).orEmpty()
}
@Suppress("DEPRECATION")
override suspend fun responseChannel(): ByteWriteChannel {
val result = ByteChannel(autoFlush = true)
if (readResponse) {
launchResponseJob(result)
}
val job = scope.reader(responseJob ?: EmptyCoroutineContext) {
val readJob = launch {
channel.copyAndClose(result, Long.MAX_VALUE)
}
configureSocketTimeoutIfNeeded(timeoutAttributes, readJob) { channel.totalBytesRead }
}
if (responseJob == null) {
responseJob = job
}
responseChannel = result
responseChannelDeferred.complete()
return job.channel
}
private fun launchResponseJob(source: ByteReadChannel) {
responseJob = async(Dispatchers.Default) {
byteContent = source.toByteArray()
}
}
override suspend fun respondOutgoingContent(content: OutgoingContent) {
super.respondOutgoingContent(content)
responseChannelDeferred.completeExceptionally(IllegalStateException("No response channel assigned"))
}
/**
* Gets a response body content channel.
*/
public fun contentChannel(): ByteReadChannel? = byteContent?.let { ByteReadChannel(it) }
internal suspend fun awaitForResponseCompletion() {
responseJob?.join()
}
// Websockets & upgrade
private val webSocketCompleted: CompletableJob = Job()
internal val webSocketEstablished: CompletableJob = Job()
override suspend fun respondUpgrade(upgrade: OutgoingContent.ProtocolUpgrade) {
upgrade.upgrade(
call.receiveChannel(),
responseChannel(),
call.application.coroutineContext,
Dispatchers.Default
).invokeOnCompletion {
webSocketCompleted.complete()
}
webSocketEstablished.complete()
}
/**
* Waits for a websocket session completion.
*/
public fun awaitWebSocket(durationMillis: Long): Unit = runBlocking {
withTimeout(durationMillis) {
responseChannelDeferred.join()
responseJob?.join()
webSocketCompleted.join()
}
Unit
}
/**
* A websocket session's channel.
*/
public fun websocketChannel(): ByteReadChannel? = responseChannel
}