Skip to content

Commit c8a188f

Browse files
committed
perf: tune video streaming path
1 parent 35e3b75 commit c8a188f

7 files changed

Lines changed: 120 additions & 4 deletions

File tree

‎src/main/kotlin/dev/typetype/server/ExtractionServiceRegistry.kt‎

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -41,6 +41,8 @@ import dev.typetype.server.services.YoutubeSessionCrypto
4141
import dev.typetype.server.services.YoutubeSessionHlsManifestService
4242
import dev.typetype.server.services.YoutubeSessionService
4343
import dev.typetype.server.services.YoutubeSessionStreamService
44+
import okhttp3.ConnectionPool
45+
import okhttp3.Dispatcher
4446
import okhttp3.OkHttpClient
4547
import java.util.concurrent.TimeUnit
4648

@@ -53,6 +55,8 @@ internal class ExtractionServiceRegistry(
5355
?.takeIf { it.length >= MIN_YOUTUBE_SESSION_SECRET_LENGTH }
5456
val httpClient = OkHttpClient()
5557
val proxyHttpClient: OkHttpClient = httpClient.newBuilder()
58+
.dispatcher(proxyDispatcher())
59+
.connectionPool(ConnectionPool(64, 5, TimeUnit.MINUTES))
5660
.connectTimeout(10, TimeUnit.SECONDS)
5761
.readTimeout(30, TimeUnit.SECONDS)
5862
.followRedirects(true)
@@ -85,7 +89,7 @@ internal class ExtractionServiceRegistry(
8589
val nicoVideoProxyService = NicoVideoProxyService()
8690
val manifestService = CachedManifestService(ManifestService(streamService), cache)
8791
val nativeManifestService = CachedNativeManifestService(NativeManifestService(), cache)
88-
val hlsManifestService = HlsManifestService(streamService, proxyHttpClient)
92+
val hlsManifestService = HlsManifestService(streamService, proxyHttpClient, cache)
8993
val youtubeSessionHlsManifestService = hlsTokenService?.let { tokenService ->
9094
youtubeSessionStreamService?.let {
9195
YoutubeSessionHlsManifestService(youtubeSessionService, it, hlsManifestService, tokenService)
@@ -95,5 +99,10 @@ internal class ExtractionServiceRegistry(
9599

96100
private companion object {
97101
const val MIN_YOUTUBE_SESSION_SECRET_LENGTH = 32
102+
103+
fun proxyDispatcher(): Dispatcher = Dispatcher().apply {
104+
maxRequests = 256
105+
maxRequestsPerHost = 64
106+
}
98107
}
99108
}

‎src/main/kotlin/dev/typetype/server/models/ProxyResponse.kt‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,4 +10,5 @@ class ProxyResponse(
1010
val acceptRanges: String?,
1111
val stream: InputStream,
1212
val close: () -> Unit,
13+
val cacheControl: String? = null,
1314
)

‎src/main/kotlin/dev/typetype/server/routes/ProxyRoutes.kt‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -28,8 +28,9 @@ fun Route.proxyRoutes(proxyService: ProxyService) {
2828
val contentType = ContentType.parse(proxy.contentType)
2929
proxy.contentRange?.let { call.response.headers.append("Content-Range", it) }
3030
proxy.acceptRanges?.let { call.response.headers.append("Accept-Ranges", it) }
31+
proxy.cacheControl?.let { call.response.headers.append("Cache-Control", it, safeOnly = false) }
3132
call.respondOutputStream(contentType, status, proxy.contentLength) {
32-
withContext(Dispatchers.IO) { proxy.stream.copyTo(this@respondOutputStream) }
33+
withContext(Dispatchers.IO) { proxy.stream.copyTo(this@respondOutputStream, PROXY_COPY_BUFFER_SIZE) }
3334
}
3435
} finally {
3536
proxy.close()
@@ -42,3 +43,5 @@ fun Route.proxyRoutes(proxyService: ProxyService) {
4243
}
4344
}
4445
}
46+
47+
private const val PROXY_COPY_BUFFER_SIZE = 64 * 1024
Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,19 @@
1+
package dev.typetype.server.services
2+
3+
import dev.typetype.server.cache.CacheService
4+
5+
internal class HlsManifestCache(private val cache: CacheService) {
6+
suspend fun get(manifestUrl: String): String? = runCatching {
7+
cache.get(key(manifestUrl))
8+
}.getOrNull()
9+
10+
suspend fun set(manifestUrl: String, manifest: String): Unit {
11+
runCatching { cache.set(key(manifestUrl), manifest, TTL_SECONDS) }
12+
}
13+
14+
private fun key(manifestUrl: String): String = PublicCacheKey.of("hls-manifest", manifestUrl)
15+
16+
private companion object {
17+
const val TTL_SECONDS = 30L
18+
}
19+
}

‎src/main/kotlin/dev/typetype/server/services/HlsManifestService.kt‎

Lines changed: 31 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,16 @@
11
package dev.typetype.server.services
22

3+
import dev.typetype.server.cache.CacheService
34
import dev.typetype.server.models.ExtractionResult
45
import dev.typetype.server.models.StreamResponse
6+
import kotlinx.coroutines.CompletableDeferred
57
import kotlinx.coroutines.Dispatchers
68
import kotlinx.coroutines.withContext
79
import okhttp3.OkHttpClient
810
import okhttp3.Request
911
import java.net.URLEncoder
1012
import java.nio.charset.StandardCharsets
13+
import java.util.concurrent.ConcurrentHashMap
1114

1215
internal fun isManifestUrl(url: String): Boolean {
1316
if (!url.startsWith("http")) return false
@@ -38,7 +41,10 @@ private fun toHlsProxyUrl(url: String): String {
3841
class HlsManifestService(
3942
private val streamService: StreamService,
4043
private val httpClient: OkHttpClient,
44+
cache: CacheService? = null,
4145
) {
46+
private val manifestCache = cache?.let(::HlsManifestCache)
47+
private val inFlight = ConcurrentHashMap<String, CompletableDeferred<ExtractionResult<String>>>()
4248

4349
suspend fun hlsManifest(url: String): ExtractionResult<String> {
4450
val manifestUrl = if (isManifestUrl(url)) {
@@ -50,7 +56,7 @@ class HlsManifestService(
5056
is ExtractionResult.Failure -> return resolved
5157
}
5258
}
53-
return fetchAndRewrite(manifestUrl)
59+
return cachedOrFetch(manifestUrl)
5460
}
5561

5662
suspend fun hlsManifestFromStreamInfo(result: ExtractionResult<StreamResponse>): ExtractionResult<String> {
@@ -59,7 +65,30 @@ class HlsManifestService(
5965
is ExtractionResult.BadRequest -> return resolved
6066
is ExtractionResult.Failure -> return resolved
6167
}
62-
return fetchAndRewrite(manifestUrl)
68+
return cachedOrFetch(manifestUrl)
69+
}
70+
71+
private suspend fun cachedOrFetch(manifestUrl: String): ExtractionResult<String> {
72+
manifestCache?.get(manifestUrl)?.let { return ExtractionResult.Success(it) }
73+
val pending = CompletableDeferred<ExtractionResult<String>>()
74+
val existing = inFlight.putIfAbsent(manifestUrl, pending)
75+
if (existing != null) return existing.await()
76+
return try {
77+
manifestCache?.get(manifestUrl)?.let {
78+
val result = ExtractionResult.Success(it)
79+
pending.complete(result)
80+
return result
81+
}
82+
val result = fetchAndRewrite(manifestUrl)
83+
if (result is ExtractionResult.Success) manifestCache?.set(manifestUrl, result.data)
84+
pending.complete(result)
85+
result
86+
} catch (error: Throwable) {
87+
pending.completeExceptionally(error)
88+
throw error
89+
} finally {
90+
inFlight.remove(manifestUrl, pending)
91+
}
6392
}
6493

6594
private suspend fun resolveHlsUrl(videoUrl: String): ExtractionResult<String> {

‎src/main/kotlin/dev/typetype/server/services/OkHttpProxyService.kt‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -57,6 +57,7 @@ class OkHttpProxyService(private val client: OkHttpClient) : ProxyService {
5757
val contentType = response.header("Content-Type") ?: "application/octet-stream"
5858
val contentRange = response.header("Content-Range")
5959
val acceptRanges = response.header("Accept-Ranges")
60+
val cacheControl = response.header("Cache-Control")
6061
val contentLength = response.header("Content-Length")?.toLongOrNull()
6162
if (isHls(contentType)) {
6263
val rewritten = if (isNicoNico(stripTrackingParams(fetchUrl))) {
@@ -71,6 +72,7 @@ class OkHttpProxyService(private val client: OkHttpClient) : ProxyService {
7172
contentLength = null,
7273
contentRange = null,
7374
acceptRanges = null,
75+
cacheControl = cacheControl,
7476
stream = ByteArrayInputStream(rewritten.toByteArray(StandardCharsets.UTF_8)),
7577
close = {},
7678
))
@@ -81,6 +83,7 @@ class OkHttpProxyService(private val client: OkHttpClient) : ProxyService {
8183
contentLength = contentLength,
8284
contentRange = contentRange,
8385
acceptRanges = acceptRanges,
86+
cacheControl = cacheControl,
8487
stream = body.byteStream(),
8588
close = response::close,
8689
))
Lines changed: 52 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,52 @@
1+
package dev.typetype.server
2+
3+
import dev.typetype.server.cache.CacheService
4+
import dev.typetype.server.models.ExtractionResult
5+
import dev.typetype.server.models.StreamResponse
6+
import dev.typetype.server.services.HlsManifestService
7+
import dev.typetype.server.services.StreamService
8+
import kotlinx.coroutines.test.runTest
9+
import okhttp3.Interceptor
10+
import okhttp3.MediaType.Companion.toMediaType
11+
import okhttp3.OkHttpClient
12+
import okhttp3.Protocol
13+
import okhttp3.Response
14+
import okhttp3.ResponseBody.Companion.toResponseBody
15+
import org.junit.jupiter.api.Assertions.assertEquals
16+
import org.junit.jupiter.api.Test
17+
18+
class HlsManifestServiceCacheTest {
19+
@Test
20+
fun `hls manifests are cached briefly by manifest url`() = runTest {
21+
var calls = 0
22+
val client = OkHttpClient.Builder().addInterceptor(Interceptor { chain ->
23+
calls += 1
24+
Response.Builder()
25+
.request(chain.request())
26+
.protocol(Protocol.HTTP_1_1)
27+
.code(200)
28+
.message("OK")
29+
.body("#EXTM3U\nsegment.ts".toResponseBody("application/vnd.apple.mpegurl".toMediaType()))
30+
.build()
31+
}).build()
32+
val service = HlsManifestService(NoopStreamService, client, InMemoryCache())
33+
val url = "https://example.com/master.m3u8"
34+
35+
service.hlsManifest(url)
36+
service.hlsManifest(url)
37+
38+
assertEquals(1, calls)
39+
}
40+
}
41+
42+
private object NoopStreamService : StreamService {
43+
override suspend fun getStreamInfo(url: String): ExtractionResult<StreamResponse> =
44+
ExtractionResult.Failure("unused")
45+
}
46+
47+
private class InMemoryCache : CacheService {
48+
private val values = mutableMapOf<String, String>()
49+
override suspend fun get(key: String): String? = values[key]
50+
override suspend fun set(key: String, value: String, ttlSeconds: Long) { values[key] = value }
51+
override suspend fun delete(key: String) { values.remove(key) }
52+
}

0 commit comments

Comments
 (0)