Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,65 @@
package app.waveflow.playback

import android.net.Uri
import androidx.media3.common.C
import androidx.media3.datasource.DataSource
import androidx.media3.datasource.DataSpec
import androidx.media3.datasource.TransferListener
import androidx.media3.datasource.cache.Cache
import androidx.media3.datasource.cache.CacheKeyFactory
import androidx.media3.datasource.cache.ContentMetadata

/**
* Retire du cache le début qu'un transcodage quitté en route y a laissé.
*
* Un transcodage en direct arrive sans longueur et refuse toute plage qui ne
* part pas du premier octet. Son début gardé en cache ne sert donc à rien : à la
* réécoute, `CacheDataSource` le relirait puis demanderait la suite par une
* plage, que le serveur refuse en 416 — une erreur que Media3 ne retente jamais.
* La piste tomberait en erreur là où le cache s'arrêtait.
*
* C'est la longueur consignée qui départage, et `CacheDataSource` la consigne
* lui-même : dès l'ouverture quand le serveur l'annonce — l'original, qui se sert
* par plages et dont le début reste utile —, à la fin du flux sinon. Une entrée
* refermée sans longueur est un transcodage abandonné avant sa fin.
*
* Ne couvre pas la **coupure réseau** en cours de transcodage : Media3 reprend
* alors à l'octet atteint, et le serveur refuse cette plage-là aussi, cache ou
* pas. C'est la relance par décalage temporel qui la couvrira.
*/
internal class IncompleteTranscodeEviction(
private val cached: DataSource,
private val cache: Cache,
private val cacheKeyFactory: CacheKeyFactory,
) : DataSource {

private var key: String? = null

override fun addTransferListener(transferListener: TransferListener) {
cached.addTransferListener(transferListener)
}

override fun open(dataSpec: DataSpec): Long {
key = cacheKeyFactory.buildCacheKey(dataSpec)
return cached.open(dataSpec)
}

override fun read(buffer: ByteArray, offset: Int, length: Int): Int = cached.read(buffer, offset, length)

override fun getUri(): Uri? = cached.uri

override fun getResponseHeaders(): Map<String, List<String>> = cached.responseHeaders

override fun close() {
try {
cached.close()
} finally {
key?.let { entree ->
if (ContentMetadata.getContentLength(cache.getContentMetadata(entree)) == C.LENGTH_UNSET.toLong()) {
cache.removeResource(entree)
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}
}
key = null
}
}
}
15 changes: 13 additions & 2 deletions app/src/main/java/app/waveflow/playback/RemoteMediaCache.kt
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import androidx.media3.datasource.DefaultDataSource
import androidx.media3.datasource.DefaultHttpDataSource
import androidx.media3.datasource.ResolvingDataSource
import androidx.media3.datasource.cache.CacheDataSource
import androidx.media3.datasource.cache.CacheKeyFactory
import androidx.media3.datasource.cache.LeastRecentlyUsedCacheEvictor
import androidx.media3.datasource.cache.SimpleCache
import kotlinx.coroutines.Dispatchers
Expand Down Expand Up @@ -49,19 +50,29 @@ class RemoteMediaCache(context: Context) : PlaybackCache {
* celle de l'élément — non l'URL de diffusion, qui change à chaque ticket
* et ne coïnciderait jamais avec elle-même ;
* - `DefaultDataSource` aiguille en amont selon le schéma : `content://` et
* `file://` partent vers les sources locales sans jamais toucher au cache.
* `file://` partent vers les sources locales sans jamais toucher au cache ;
* - au-dessus du cache, [IncompleteTranscodeEviction] retire le début d'un
* transcodage quitté en route, que le serveur ne laisserait pas compléter.
*/
fun dataSourceFactory(resolver: ResolvingDataSource.Resolver): DataSource.Factory {
val resolving = ResolvingDataSource.Factory(DefaultHttpDataSource.Factory(), resolver)
// Une seule fabrique de clés pour les deux : l'éviction doit viser
// l'entrée même que le cache vient d'écrire.
val cles = CacheKeyFactory.DEFAULT

val cached = CacheDataSource.Factory()
.setCache(cache)
.setCacheKeyFactory(cles)
.setUpstreamDataSourceFactory(resolving)
// Un cache illisible doit dégrader vers le réseau, pas interrompre
// la lecture.
.setFlags(CacheDataSource.FLAG_IGNORE_CACHE_ON_ERROR)

return DefaultDataSource.Factory(appContext, cached)
val sansDebutOrphelin = DataSource.Factory {
IncompleteTranscodeEviction(cached.createDataSource(), cache, cles)
}

return DefaultDataSource.Factory(appContext, sansDebutOrphelin)
}

override val maxBytes: Long = MAX_BYTES
Expand Down
140 changes: 133 additions & 7 deletions app/src/test/java/app/waveflow/playback/RemoteMediaCacheTest.kt
Original file line number Diff line number Diff line change
Expand Up @@ -5,13 +5,19 @@ import android.content.ContentValues
import android.net.Uri
import android.os.ParcelFileDescriptor
import androidx.core.net.toUri
import androidx.media3.common.C
import androidx.media3.common.MediaItem
import androidx.media3.datasource.DataSource
import androidx.media3.datasource.DataSpec
import androidx.media3.datasource.ResolvingDataSource
import androidx.test.core.app.ApplicationProvider
import app.waveflow.model.StreamRendering
import kotlinx.coroutines.test.runTest
import okhttp3.mockwebserver.Dispatcher
import okhttp3.mockwebserver.MockResponse
import okhttp3.mockwebserver.MockWebServer
import okhttp3.mockwebserver.RecordedRequest
import okio.Buffer
import org.junit.After
import org.junit.Assert.assertArrayEquals
import org.junit.Assert.assertEquals
Expand Down Expand Up @@ -57,11 +63,16 @@ class RemoteMediaCacheTest {
server.shutdown()
}

/** Un résolveur qui rend une URL différente à chaque appel, comme le vrai. */
/**
* Un résolveur qui rend une URL différente à chaque appel, comme le vrai.
*
* Le rendu du marqueur suit jusqu'au serveur, comme dans [RemoteStreamResolver].
*/
private val resolver = ResolvingDataSource.Resolver { dataSpec ->
val trackId = trackIdOfRemoteUri(dataSpec.uri) ?: return@Resolver dataSpec
ticketsDemandes++
dataSpec.withUri(server.url("/api/v2/stream/ticket-$ticketsDemandes-$trackId").toString().toUri())
val rendu = dataSpec.uri.query?.let { "?$it" }.orEmpty()
dataSpec.withUri(server.url("/api/v2/stream/ticket-$ticketsDemandes-$trackId$rendu").toString().toUri())
}

private fun lire(source: DataSource, spec: DataSpec): ByteArray {
Expand All @@ -71,7 +82,7 @@ class RemoteMediaCacheTest {
val sortie = java.io.ByteArrayOutputStream()
while (true) {
val lus = source.read(tampon, 0, tampon.size)
if (lus == androidx.media3.common.C.RESULT_END_OF_INPUT) break
if (lus == C.RESULT_END_OF_INPUT) break
sortie.write(tampon, 0, lus)
}
sortie.toByteArray()
Expand All @@ -80,14 +91,73 @@ class RemoteMediaCacheTest {
}
}

/** Lit [octets] puis referme, comme le lecteur qui passe au morceau suivant. */
private fun lireLeDebut(source: DataSource, spec: DataSpec, octets: Int) {
source.open(spec)
try {
val tampon = ByteArray(octets)
var lus = 0
while (lus < octets) {
val n = source.read(tampon, lus, octets - lus)
if (n == C.RESULT_END_OF_INPUT) break
lus += n
}
} finally {
source.close()
}
}

/** Le marqueur et la clé qu'une piste distante porte réellement dans la file. */
private fun specDe(
trackId: String,
format: String = DEFAULT_FORMAT,
bitrate: Int? = null,
) = DataSpec.Builder()
.setUri("waveflow://track/$trackId".toUri())
.setKey(cacheKeyOf(trackId, format, bitrate))
.build()
): DataSpec {
val piste = MediaItem.Builder()
.setUri("waveflow://track/$trackId".toUri())
.build()
.withRendering(StreamRendering(format, bitrate))
val configuration = checkNotNull(piste.localConfiguration)
return DataSpec.Builder()
.setUri(configuration.uri)
.setKey(configuration.customCacheKey)
.build()
}

/**
* Répond comme `waveflow-server` (`src/media.rs`) : l'original se sert par
* plages ; un transcodage en direct arrive par morceaux, sans longueur, et
* refuse toute plage qui ne part pas du premier octet.
*/
private fun servirCommeLeServeur(contenu: ByteArray) {
server.dispatcher = object : Dispatcher() {
override fun dispatch(request: RecordedRequest): MockResponse {
val debut = request.getHeader("Range")
?.removePrefix("bytes=")
?.substringBefore('-')
?.toInt()
?: 0
val transcode = request.requestUrl?.queryParameter("format") != null
return when {
transcode && debut > 0 -> MockResponse()
.setResponseCode(416)
.setHeader("Content-Range", "bytes */0")
.setHeader("Accept-Ranges", "none")

transcode -> MockResponse()
.setHeader("Accept-Ranges", "none")
.setChunkedBody(Buffer().write(contenu), 4096)

debut > 0 -> MockResponse()
.setResponseCode(206)
.setHeader("Content-Range", "bytes $debut-${contenu.size - 1}/${contenu.size}")
.setBody(Buffer().write(contenu, debut, contenu.size - debut))

else -> MockResponse().setBody(Buffer().write(contenu))
}
}
}
}

@Test
fun `une seconde lecture ne redemande ni octets ni ticket`() {
Expand Down Expand Up @@ -168,6 +238,57 @@ class RemoteMediaCacheTest {
assertEquals("transcode", String(transcode))
}

@Test
fun `un transcodage quitte en route se relit en entier`() {
// On passe au morceau suivant avant la fin : le cache garde le début.
// Le relire puis demander la suite, c'est demander une plage à un
// transcodage en direct, qui la refuse — et Media3 ne retente pas un
// 416. La piste tombait en erreur là où le cache s'arrêtait.
val contenu = octetsAudio()
servirCommeLeServeur(contenu)
val factory = mediaCache.dataSourceFactory(resolver)

lireLeDebut(factory.createDataSource(), specDe("piste-1", "opus", 96), OCTETS_ECOUTES)
val relu = lire(factory.createDataSource(), specDe("piste-1", "opus", 96))

assertArrayEquals(contenu, relu)
}

@Test
fun `un transcodage lu jusqu'au bout reste en cache`() {
// Le pendant du précédent. Sa longueur n'arrive qu'avec la fin du flux,
// et c'est elle qui distingue le morceau entier d'un début abandonné.
val contenu = octetsAudio()
servirCommeLeServeur(contenu)
val factory = mediaCache.dataSourceFactory(resolver)

lire(factory.createDataSource(), specDe("piste-1", "opus", 96))
val relu = lire(factory.createDataSource(), specDe("piste-1", "opus", 96))

assertArrayEquals(contenu, relu)
assertEquals("une seule requête réseau", 1, server.requestCount)
}

@Test
fun `un original quitte en route reprend la ou le cache s'arrete`() {
// L'original, lui, se sert par plages : son début en cache reste utile,
// et le jeter ferait retélécharger ce qu'on a déjà.
val contenu = octetsAudio()
servirCommeLeServeur(contenu)
val factory = mediaCache.dataSourceFactory(resolver)

lireLeDebut(factory.createDataSource(), specDe("piste-1"), OCTETS_ECOUTES)
val relu = lire(factory.createDataSource(), specDe("piste-1"))

assertArrayEquals(contenu, relu)
server.takeRequest()
assertEquals(
"seule la suite est redemandée",
"bytes=$OCTETS_ECOUTES-${contenu.size - 1}",
server.takeRequest().getHeader("Range"),
)
}

@Test
fun `un fichier local ne passe pas par le cache`() {
// Il est déjà sur le disque : le recopier doublerait sa place.
Expand Down Expand Up @@ -250,6 +371,11 @@ class RemoteMediaCacheTest {

private const val AUTORITE = "app.waveflow.test.audio"

/** Assez pour qu'un début lu laisse une vraie suite à demander. */
private const val OCTETS_ECOUTES = 8_000

private fun octetsAudio() = ByteArray(64_000) { (it % 251).toByte() }

/**
* Sert un fichier temporaire derrière une URI `content://`.
*
Expand Down