Merge pull request #4172 from vitorpamplona/claude/morganite-blossom-video-issue-e239x2

Local Blossom cache: fix video playback, harden the bridge, route it only at the HTTP layer
This commit is contained in:
Vitor Pamplona
2026-09-23 14:05:09 -04:00
committed by GitHub
26 changed files with 999 additions and 718 deletions
@@ -50,6 +50,7 @@ import com.vitorpamplona.amethyst.commons.service.http.BlossomReadAuthTokenProvi
import com.vitorpamplona.amethyst.commons.service.http.DualHttpClientManager
import com.vitorpamplona.amethyst.commons.service.http.DualHttpClientManagerForRelays
import com.vitorpamplona.amethyst.commons.service.http.EncryptionKeyCache
import com.vitorpamplona.amethyst.commons.service.http.LocalBlossomMediaCallFactory
import com.vitorpamplona.amethyst.commons.service.http.OnionLocationCache
import com.vitorpamplona.amethyst.commons.service.lnurl.OkHttpLnurlEndpointResolver
import com.vitorpamplona.amethyst.commons.service.pow.PoWPolicy
@@ -458,13 +459,14 @@ class AppModules(
scope = applicationIOScope,
dns = surgeDns,
// Transparently rewrites sha256-keyed HTTP requests to the local
// Blossom cache when the master toggle is on, the profile-pictures-only
// restriction is off, and the probe sees 127.0.0.1:24242 as available.
shouldBridgeBlossomCache = {
// Blossom cache when the master toggle is on, the probe sees
// 127.0.0.1:24242 as available, and either the profile-pictures-only
// restriction is off or the request is a profile picture.
shouldBridgeBlossomCache = { profilePicture ->
val settings = sessionManager.loggedInAccount()?.settings
val master = settings?.useLocalBlossomCache?.value ?: false
val profileOnly = settings?.localBlossomCacheProfilePicturesOnly?.value ?: false
master && !profileOnly && localBlossomCacheProbe.available.value
master && (profilePicture || !profileOnly) && localBlossomCacheProbe.available.value
},
onionCache = onionLocationCache,
usageInterceptor = httpUsageInterceptor,
@@ -477,6 +479,7 @@ class AppModules(
cachedHeaderProvider = blossomReadAuthTokens::cachedHeader,
onAuthRequired = blossomReadAuthTokens::warm,
),
onLocalBlossomCacheUnreachable = { localBlossomCacheProbe.markUnavailable() },
)
// Offers easy methods to know when connections are happening through Tor or not
@@ -1035,14 +1038,16 @@ class AppModules(
}
},
httpClientBuilder = roleBasedHttpClientBuilder,
// Same gate as the interceptor: the profile-pictures-only restriction keeps note
// media (including native `blossom:` URIs) off the local cache.
useLocalBlossomCache = {
sessionManager
.loggedInAccount()
?.settings
?.useLocalBlossomCache
?.value ?: false
val settings = sessionManager.loggedInAccount()?.settings
val master = settings?.useLocalBlossomCache?.value ?: false
val profileOnly = settings?.localBlossomCacheProfilePicturesOnly?.value ?: false
master && !profileOnly
},
localCacheProbe = localBlossomCacheProbe,
scope = applicationIOScope,
)
}
@@ -1148,7 +1153,8 @@ class AppModules(
blossomServerResolver = { blossomResolver },
// Through the role builder (not raw getHttpClient) so Coil's image
// traffic carries the "image" ledger tag. Same Tor decision inside.
callFactory = { roleBasedHttpClientBuilder.okHttpClientForImage(it) },
// Marked as media so the local Blossom cache may serve these downloads.
callFactory = { LocalBlossomMediaCallFactory(roleBasedHttpClientBuilder.okHttpClientForImage(it)) },
thumbnailCache = thumbnailDiskCache,
backgroundScope = applicationIOScope,
readAuth = blossomReadAuthTokens,
@@ -1341,8 +1347,7 @@ class AppModules(
state.account.settings.localBlossomCacheProfilePicturesOnly
.drop(1),
).collect {
blossomResolver.uriToUrlCache.evictAll()
blossomResolver.blossomHitCache.cache.evictAll()
blossomResolver.clearCaches()
localBlossomCacheProbe.invalidate()
// Re-probe immediately so enabling the feature activates it
// this session. Otherwise `available` only advances when a
@@ -1358,8 +1363,14 @@ class AppModules(
}
applicationIOScope.launch {
localBlossomCacheProbe.available.drop(1).collect {
blossomResolver.uriToUrlCache.evictAll()
blossomResolver.blossomHitCache.cache.evictAll()
blossomResolver.clearCaches()
}
}
// A resolution to the local cache made while a media type was not Tor-routed must not
// survive the user switching it to Tor: the cache would keep fetching it outside Tor.
applicationIOScope.launch {
torPrefs.value.propertyWatchFlow.drop(1).collect {
blossomResolver.clearCaches()
}
}
// Warm the local-cache probe so the very first image load doesn't pay
@@ -1367,6 +1378,22 @@ class AppModules(
applicationIOScope.launch {
localBlossomCacheProbe.isAvailable()
}
// Re-checks the local cache about once a minute while the feature is on, so the bridge
// comes back when the cache app is restarted and turns off when it is closed, even while
// no `blossom:` URI is being resolved. A loopback HEAD never wakes the radio.
applicationIOScope.launch {
while (true) {
delay(LOCAL_BLOSSOM_CACHE_RECHECK_MS)
if (sessionManager
.loggedInAccount()
?.settings
?.useLocalBlossomCache
?.value == true
) {
localBlossomCacheProbe.isAvailable()
}
}
}
// Warms the video cache off the main thread. SimpleCache's constructor opens a SQLite
// index over StandaloneDatabaseProvider and walks every cached span on disk — up to a
@@ -1505,5 +1532,8 @@ class AppModules(
* check and burn CPU on a heap that has nothing left to give.
*/
private const val MIN_RECLAIM_INTERVAL_MS = 120_000L
/** How often the local Blossom cache is re-probed while the feature is enabled. */
private const val LOCAL_BLOSSOM_CACHE_RECHECK_MS = 60_000L
}
}
@@ -36,6 +36,7 @@ import com.vitorpamplona.amethyst.commons.service.http.BlossomReadAuthTokenProvi
import com.vitorpamplona.amethyst.commons.service.image.readAuthAware
import com.vitorpamplona.amethyst.commons.service.image.withAuthHeader
import com.vitorpamplona.amethyst.service.uploads.blossom.bud10.BlossomServerResolver
import com.vitorpamplona.quartz.nipB7Blossom.BlossomUri
import com.vitorpamplona.quartz.utils.startsWithIgnoreCase
import okhttp3.Call
import kotlin.coroutines.cancellation.CancellationException
@@ -44,18 +45,49 @@ import kotlin.coroutines.cancellation.CancellationException
class BlossomFetcher(
private val options: Options,
private val data: Uri,
private val imageLoader: ImageLoader,
private val blossomServerResolver: () -> BlossomServerResolver,
private val networkFetcher: (url: String) -> Fetcher,
private val networkFetcher: (url: String, diskCacheKey: String) -> Fetcher,
) : Fetcher {
override suspend fun fetch(): FetchResult? =
try {
val urlResult = blossomServerResolver().findServers(data.toString())
networkFetcher(urlResult?.serverUrl ?: data.toString()).fetch()
val uri = data.toString()
val key = options.diskCacheKey ?: diskCacheKey(uri)
// A blob already on disk is served without resolving: resolution may HEAD
// servers and wait on relays, and fails outright offline, where the bytes are
// still sitting in the cache. NetworkFetcher reads the snapshot under [key]
// before it ever looks at the url.
val onDisk = options.diskCachePolicy.readEnabled && imageLoader.diskCache?.openSnapshot(key)?.use { true } ?: false
val fromDisk =
if (onDisk) {
// Evicted between the check and the read: the fetcher would try the
// network with the unresolvable `blossom:` uri, so resolve after all.
try {
networkFetcher(uri, key).fetch()
} catch (e: Exception) {
if (e is CancellationException) throw e
null
}
} else {
null
}
fromDisk ?: networkFetcher(blossomServerResolver().findServers(uri)?.serverUrl ?: uri, key).fetch()
} catch (e: Exception) {
if (e is CancellationException) throw e
null
}
companion object {
/**
* Blossom blobs are content-addressed, so the sha256 names the bytes wherever they
* were served from. Keying on it (instead of the resolved server URL) finds the
* blob again after a restart, a probe flip or a different server answering.
*/
fun diskCacheKey(uri: String): String = BlossomUri.parse(uri)?.let { "blossom:${it.sha256}" } ?: uri
}
@OptIn(ExperimentalCoilApi::class)
class Factory(
val blossomServerResolver: () -> BlossomServerResolver,
@@ -78,20 +110,19 @@ class BlossomFetcher(
if (!isApplicable(data)) return null
// Wrapped per resolved url (not per Factory) because the server the
// blob actually lives on is only known once the resolver has run.
return BlossomFetcher(options, data, blossomServerResolver) { url ->
return BlossomFetcher(options, data, imageLoader, blossomServerResolver) { url, diskCacheKey ->
val keyedOptions = options.copy(diskCacheKey = diskCacheKey)
readAuthAware(url, readAuth) { authHeader ->
NetworkFetcher(
url = url,
options = options.withAuthHeader(authHeader),
options = keyedOptions.withAuthHeader(authHeader),
networkClient = lazy { networkClient(url).asNetworkClient() },
diskCache = lazy { imageLoader.diskCache },
cacheStrategy = cacheStrategyLazy,
connectivityChecker = lazy { connectivityCheckerLazy.get(options.context) },
concurrentRequestStrategy = concurrentRequestStrategyLazy,
)
// Keyed on the resolved server url, which is what the NetworkFetcher above
// caches under -- not on the `blossom:` uri the request came in as.
}.onSystemFileSystem(options.diskCacheKey ?: url)
}.onSystemFileSystem(diskCacheKey)
}
}
@@ -87,9 +87,10 @@ class ImageLoaderSetup {
// coordinates through a map of in-flight fetches that it owns, so it only
// works when every fetcher shares the same instance -- a fresh one per
// request can never see anybody else's fetch and the de-dupe silently
// no-ops. Shared across all three network-backed factories so a feed
// image and the same blob reached through `blossom:` (or a profile
// picture) still collapse onto one download.
// no-ops. Shared across all three network-backed factories so the same
// key requested through any of them (e.g. a feed image and a profile
// picture with the same URL) collapses onto one download. `blossom:`
// blobs are keyed by their sha256 instead (see BlossomFetcher.diskCacheKey).
val concurrentRequests = DeDupeConcurrentRequestStrategy()
SingletonImageLoader.setUnsafe(
@@ -36,6 +36,8 @@ import coil3.network.NetworkFetcher
import coil3.network.okhttp.asNetworkClient
import coil3.request.Options
import com.vitorpamplona.amethyst.commons.service.http.BlossomReadAuthTokenProvider
import com.vitorpamplona.amethyst.commons.service.http.LocalBlossomCacheRedirectInterceptor.Companion.PROFILE_PICTURE
import com.vitorpamplona.amethyst.commons.service.http.LocalBlossomMediaCallFactory
import com.vitorpamplona.amethyst.commons.service.image.readAuthAware
import com.vitorpamplona.amethyst.commons.service.image.withAuthHeader
import com.vitorpamplona.amethyst.commons.ui.components.ProfilePictureUrl
@@ -121,7 +123,8 @@ class ProfilePictureFetcher(
NetworkFetcher(
url = data.url,
options = options.withAuthHeader(authHeader),
networkClient = lazy { networkClient(data.url).asNetworkClient() },
// Marked so the local Blossom cache bridge can honour "profile pictures only".
networkClient = lazy { LocalBlossomMediaCallFactory(networkClient(data.url), PROFILE_PICTURE).asNetworkClient() },
diskCache = diskCacheLazy,
cacheStrategy = cacheStrategyLazy,
connectivityChecker = lazy { connectivityCheckerLazy.get(options.context) },
@@ -39,6 +39,7 @@ import androidx.media3.session.MediaSession
import androidx.media3.session.MediaSessionService
import com.vitorpamplona.amethyst.Amethyst
import com.vitorpamplona.amethyst.commons.service.http.DynamicCallFactory
import com.vitorpamplona.amethyst.commons.service.http.LocalBlossomCacheRedirectInterceptor
import com.vitorpamplona.amethyst.service.playback.diskCache.VideoCache
import com.vitorpamplona.amethyst.service.playback.pip.BackgroundMedia
import com.vitorpamplona.amethyst.service.playback.playerPool.ExoPlayerBuilder
@@ -59,7 +60,11 @@ class PlaybackService : MediaSessionService() {
okHttpClient: DynamicCallFactory,
blossomServerResolver: BlossomServerResolver,
): MediaSessionPool {
val dataSourceFactory = OkHttpDataSource.Factory(okHttpClient)
// Marked as media so the local Blossom cache may serve these downloads.
val dataSourceFactory =
OkHttpDataSource
.Factory(okHttpClient)
.setDefaultRequestProperties(mapOf(LocalBlossomCacheRedirectInterceptor.MEDIA_HEADER to LocalBlossomCacheRedirectInterceptor.MEDIA))
val resolvingDataSourceFactory: DataSource.Factory =
ResolvingDataSource.Factory(
@@ -23,12 +23,14 @@ package com.vitorpamplona.amethyst.service.resourceusage
import android.os.SystemClock
import okhttp3.Interceptor
import okhttp3.OkHttpClient
import okhttp3.Request
import okhttp3.Response
import okhttp3.ResponseBody
import okio.Buffer
import okio.ForwardingSource
import okio.Source
import okio.buffer
import java.io.IOException
import java.util.concurrent.ConcurrentHashMap
import java.util.concurrent.atomic.AtomicBoolean
@@ -97,24 +99,29 @@ class UsageCountingInterceptor(
override fun intercept(chain: Interceptor.Chain): Response {
val request = chain.request()
// Loopback traffic (LocalBlossomCacheRedirectInterceptor rewrites cache
// hits to 127.0.0.1) never touches the radio: counting it would inflate
// the network numbers — and could trip the background-data alert — for
// exactly the users who cache aggressively to SAVE data.
// Loopback traffic never touches the radio: counting it would inflate the
// network numbers — and could trip the background-data alert — for exactly
// the users who cache aggressively to SAVE data. This interceptor is the
// outermost one, so it sees URLs BEFORE LocalBlossomCacheRedirectInterceptor
// rewrites them to 127.0.0.1; the host that was actually contacted is only
// known from the response's request.
if (isLoopback(request.url.host)) return chain.proceed(request)
val role = request.tag(UsageRoleTag::class.java)?.role ?: defaultRole
accountant.add(UsageKeys.netReqs(role, isMobile(), isForeground()), 1)
bursts?.onHttpActivity()
val requestBytes = request.body?.contentLength()?.coerceAtLeast(0L) ?: 0L
if (requestBytes > 0) {
accountant.add(UsageKeys.net(role, isMobile(), isForeground(), received = false), requestBytes)
}
val startedAtMs = nowMs()
val response = chain.proceed(request)
val response =
try {
chain.proceed(request)
} catch (e: IOException) {
countRequest(request, role)
throw e
}
if (isLoopback(response.request.url.host)) return response
countRequest(request, role)
return response
.newBuilder()
.body(
@@ -131,6 +138,19 @@ class UsageCountingInterceptor(
).build()
}
private fun countRequest(
request: Request,
role: String,
) {
accountant.add(UsageKeys.netReqs(role, isMobile(), isForeground()), 1)
bursts?.onHttpActivity()
val requestBytes = request.body?.contentLength()?.coerceAtLeast(0L) ?: 0L
if (requestBytes > 0) {
accountant.add(UsageKeys.net(role, isMobile(), isForeground(), received = false), requestBytes)
}
}
companion object {
/** OkHttp reports IPv6 hosts unbracketed ("::1"); keep the bracketed form defensively. */
fun isLoopback(host: String): Boolean = host.startsWith("127.") || host == "localhost" || host == "::1" || host == "[::1]"
@@ -28,17 +28,23 @@ import com.vitorpamplona.quartz.nip01Core.core.HexKey
import com.vitorpamplona.quartz.nip01Core.core.isValid
import com.vitorpamplona.quartz.nipB7Blossom.BlossomServersEvent
import com.vitorpamplona.quartz.nipB7Blossom.BlossomUri
import com.vitorpamplona.quartz.utils.TimeUtils
import com.vitorpamplona.quartz.utils.firstNotNullOrNullAsync
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Deferred
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.ExperimentalCoroutinesApi
import kotlinx.coroutines.IO
import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.async
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.first
import kotlinx.coroutines.flow.merge
import kotlinx.coroutines.flow.transformLatest
import kotlinx.coroutines.withContext
import kotlinx.coroutines.withTimeoutOrNull
import okhttp3.OkHttpClient
import java.util.concurrent.ConcurrentHashMap
import java.util.concurrent.atomic.AtomicInteger
class BlossomServerResolver(
val loggedInUsers: () -> List<HexKey>,
@@ -46,10 +52,20 @@ class BlossomServerResolver(
val httpClientBuilder: IRoleBasedHttpClientBuilder,
val useLocalBlossomCache: () -> Boolean = { false },
val localCacheProbe: LocalBlossomCacheProbe? = null,
private val scope: CoroutineScope = CoroutineScope(SupervisorJob() + Dispatchers.IO),
) {
val blossomHitCache: ServerHeadCache = ServerHeadCache()
val uriToUrlCache = LruCache<String, BlossomUriServer>(200)
// Unresolvable URIs, with when they failed. Without this every retry — each ExoPlayer
// load attempt blocks a loader thread on it — paid the full resolution timeout again.
private val missCache = LruCache<String, Long>(200)
// One resolution per URI at a time: the feed preview, Coil and the player asking for the
// same `blossom:` URI together share it instead of each firing its own HEADs and relay
// subscriptions. Runs in [scope] so a caller that goes away doesn't cancel it for the rest.
private val inFlight = ConcurrentHashMap<String, Deferred<BlossomUriServer?>>()
class BlossomUriServer(
val uri: BlossomUri,
val serverUrl: String,
@@ -57,41 +73,71 @@ class BlossomServerResolver(
fun cachedFindServer(uriStr: String): BlossomUriServer? = uriToUrlCache[uriStr]
// Bumped by [clearCaches]. A resolution that started before a clear (e.g. it picked the
// local cache just before the cache went down) must not write its stale answer back.
private val generation = AtomicInteger(0)
/** Forgets every resolution, e.g. when the local cache comes up or goes down. */
fun clearCaches() {
generation.incrementAndGet()
uriToUrlCache.evictAll()
missCache.evictAll()
blossomHitCache.cache.evictAll()
}
suspend fun findServers(uriStr: String): BlossomUriServer? {
uriToUrlCache[uriStr]?.let { return it }
missCache[uriStr]?.let { failedAt ->
if (TimeUtils.nowMillis() - failedAt < MISS_TTL_MS) return null
}
// Confined to Dispatchers.IO: this is reached from Compose
// `produceState`/`LaunchedEffect` (RichTextViewer, MarmotGroupIconDisplay),
// which run on the main dispatcher. The pre-suspension work here —
// BlossomUri parsing, LruCache lookups, the local-cache probe's client
// build, and the server-list flow setup — must stay off the UI thread.
val result =
withContext(Dispatchers.IO) {
withTimeoutOrNull(10000) {
findServersInner(uriStr)
val resolution =
inFlight.computeIfAbsent(uriStr) {
// Confined to Dispatchers.IO: this is reached from Compose
// `produceState`/`LaunchedEffect` (RichTextViewer, MarmotGroupIconDisplay),
// which run on the main dispatcher. The pre-suspension work here —
// BlossomUri parsing, LruCache lookups, the local-cache probe's client
// build, and the server-list flow setup — must stay off the UI thread.
scope.async(Dispatchers.IO) {
val startedAt = generation.get()
try {
val result = withTimeoutOrNull(RESOLVE_TIMEOUT_MS) { findServersInner(uriStr) }
// Cleared while resolving: answer the waiting callers, remember nothing.
if (startedAt == generation.get()) {
if (result != null) {
uriToUrlCache.put(uriStr, result)
} else {
missCache.put(uriStr, TimeUtils.nowMillis())
}
}
result
} finally {
inFlight.remove(uriStr)
}
}
}
if (result != null) {
uriToUrlCache.put(uriStr, result)
}
return result
return resolution.await()
}
@OptIn(ExperimentalCoroutinesApi::class)
suspend fun findServersInner(uriStr: String): BlossomUriServer? {
val uri = BlossomUri.parse(uriStr) ?: return null
if (useLocalBlossomCache() && localCacheProbe?.isAvailable() == true) {
val expectedMimeType = mimeTypeMap[uri.extension]
// Same rule as the HTTP-level bridge, which Tor-routed clients don't carry: media the
// user sends through Tor must not be handed to the local cache, which would fetch it
// from the origin outside Tor.
if (useLocalBlossomCache() && !isTorRouted(uri, expectedMimeType) && localCacheProbe?.isAvailable() == true) {
return BlossomUriServer(uri, uri.toLocalCacheUrl(LocalBlossomCacheProbe.LOCAL_CACHE_BASE))
}
val expectedMimeType = mimeTypeMap[uri.extension]
val filename = uri.filename()
if (uri.servers.isNotEmpty()) {
val workingUrl = firstWorkingUrl(uri.servers, filename, expectedMimeType, uri.size)
// Bounded well under RESOLVE_TIMEOUT_MS so a hint server that drops packets
// leaves time for the author's own server list below.
val workingUrl = firstWorkingUrl(uri.servers, filename, expectedMimeType, uri.size, XS_TIMEOUT_MS)
if (workingUrl != null) {
return BlossomUriServer(uri, workingUrl)
}
@@ -135,13 +181,20 @@ class BlossomServerResolver(
filename: String,
expectedMimeType: String?,
expectedSize: Long?,
timeoutMs: Long = RESOLVE_TIMEOUT_MS,
): String? =
firstNotNullOrNullAsync(servers, 10000) {
firstNotNullOrNullAsync(servers, timeoutMs) {
blossomHitCache.urlIfServerHasFile(it, filename, expectedMimeType, expectedSize) { url ->
client(url, expectedMimeType)
}
}
/** Whether this blob's media type would be fetched through Tor from a regular (clearnet) server. */
private fun isTorRouted(
uri: BlossomUri,
mimeType: String?,
): Boolean = client(uri.toServerUrl() ?: "https://blossom.invalid/${uri.filename()}", mimeType).proxy != null
fun client(
url: String,
mimeType: String?,
@@ -153,9 +206,13 @@ class BlossomServerResolver(
else -> httpClientBuilder.okHttpClientForPreview(url)
}
fun canResolve(scheme: String) = scheme == SCHEME
fun canResolve(scheme: String) = scheme.equals(SCHEME, ignoreCase = true)
companion object {
const val SCHEME = "blossom"
private const val RESOLVE_TIMEOUT_MS = 10_000L
private const val XS_TIMEOUT_MS = 4_000L
private const val MISS_TTL_MS = 30_000L
}
}
@@ -32,6 +32,7 @@ import kotlinx.coroutines.withContext
import okhttp3.Request
import okhttp3.coroutines.executeAsync
import java.util.concurrent.TimeUnit
import java.util.concurrent.atomic.AtomicInteger
/**
* Discovers a local Blossom cache running on `http://127.0.0.1:24242` per
@@ -48,6 +49,9 @@ class LocalBlossomCacheProbe(
@Volatile
private var cachedAtMs: Long = 0L
// Bumped by [markUnavailable]; lets a probe that raced it discard its stale positive result.
private val unavailableMarks = AtomicInteger(0)
private val _available = MutableStateFlow(false)
val available: StateFlow<Boolean> = _available
@@ -66,7 +70,11 @@ class LocalBlossomCacheProbe(
return@withLock _available.value
}
val startedAt = unavailableMarks.get()
val newResult = probe()
// A refused connection reported while this probe was in flight is newer evidence
// than a HEAD that may have succeeded just before the cache went away.
if (newResult && startedAt != unavailableMarks.get()) return@withLock false
_available.value = newResult
cachedAtMs = currentTimeMs()
newResult
@@ -80,6 +88,17 @@ class LocalBlossomCacheProbe(
cachedAtMs = 0L
}
/**
* Records that the cache just refused a connection. Flipping [available] right away turns
* the bridge off everywhere it is read, instead of letting every sha256 URL keep failing
* against the dead loopback port until the positive TTL runs out and someone re-probes.
*/
fun markUnavailable() {
unavailableMarks.incrementAndGet()
cachedAtMs = currentTimeMs()
_available.value = false
}
// Confined to Dispatchers.IO because callers reach this through suspend
// resolvers invoked from Compose `LaunchedEffect`/`produceState`, which run
// on the main dispatcher: building the OkHttp client and issuing the HEAD
@@ -22,7 +22,7 @@ package com.vitorpamplona.amethyst.service.uploads.blossom.bud10
import androidx.collection.LruCache
import kotlinx.coroutines.CancellationException
import okhttp3.MediaType.Companion.toMediaType
import okhttp3.MediaType.Companion.toMediaTypeOrNull
import okhttp3.OkHttpClient
import okhttp3.Request
import okhttp3.coroutines.executeAsync
@@ -56,13 +56,18 @@ class ServerHeadCache {
client(url).newCall(request).executeAsync().use { response ->
if (!response.isSuccessful) {
cache.put(url, HasFile.NoFile)
// Only a definitive answer is remembered. A 5xx, 429 or auth challenge is
// transient: caching it would hide this server for the app's lifetime.
if (response.code == 404 || response.code == 410) {
cache.put(url, HasFile.NoFile)
}
return HasFile.NoFile
}
// Retrieve the "Content-Length" header
val contentLength = response.header("Content-Length")?.toLongOrNull()
val mimeType = response.header("Content-Type")?.toMediaType()?.toString()
// type/subtype only: a `; charset=binary` parameter must not fail the match below.
val mimeType = response.header("Content-Type")?.toMediaTypeOrNull()?.let { "${it.type}/${it.subtype}" }
return if (contentLength != null && mimeType != null) {
val result = HasFile.TypeAndSize(mimeType, contentLength)
@@ -75,7 +80,8 @@ class ServerHeadCache {
}
} catch (e: Exception) {
if (e is CancellationException) throw e
cache.put(url, HasFile.NoFile)
// Offline, DNS failure, timeout: says nothing about the server having the file,
// so it is not cached — a brief outage must not hide the blob until eviction.
return HasFile.NoFile
}
}
@@ -104,7 +110,7 @@ class ServerHeadCache {
if (result.size == expectedSize) {
return url
}
if (expectedSize == null && result.size > 0 && result.mimeType == expectedMimeType) {
if (expectedSize == null && result.size > 0 && result.mimeType.equals(expectedMimeType, ignoreCase = true)) {
return url
}
}
@@ -212,7 +212,7 @@ fun MediaCacheSection(accountViewModel: AccountViewModel) {
*/
@Composable
private fun CacheDetectionChip(accountViewModel: AccountViewModel) {
val probeAvailable by accountViewModel.useLocalBlossomBridgeForProfilePics
val probeAvailable by accountViewModel.localBlossomCacheDetected
.collectAsStateWithLifecycle()
val color = if (probeAvailable) MaterialTheme.colorScheme.allGoodColor else MaterialTheme.colorScheme.grayText
@@ -39,48 +39,21 @@ import androidx.compose.ui.graphics.drawscope.DrawScope
import androidx.compose.ui.graphics.vector.rememberVectorPainter
import androidx.compose.ui.layout.ContentScale
import androidx.compose.ui.platform.LocalContext
import androidx.lifecycle.compose.collectAsStateWithLifecycle
import coil3.asDrawable
import coil3.compose.AsyncImagePainter
import coil3.compose.SubcomposeAsyncImage
import coil3.compose.SubcomposeAsyncImageContent
import com.vitorpamplona.amethyst.Amethyst
import coil3.network.NetworkHeaders
import coil3.network.httpHeaders
import coil3.request.ImageRequest
import com.vitorpamplona.amethyst.commons.icons.symbols.MaterialSymbols
import com.vitorpamplona.amethyst.commons.icons.symbols.rememberMaterialSymbolPainter
import com.vitorpamplona.amethyst.commons.richtext.bridgeProfilePictureUrl
import com.vitorpamplona.amethyst.commons.robohash.CachedRobohash
import com.vitorpamplona.amethyst.commons.service.http.LocalBlossomCacheRedirectInterceptor
import com.vitorpamplona.amethyst.commons.ui.components.ProfilePictureUrl
import com.vitorpamplona.amethyst.commons.ui.components.forwardingPainter
import com.vitorpamplona.amethyst.ui.screen.AccountState
import com.vitorpamplona.amethyst.ui.theme.isLight
import com.vitorpamplona.amethyst.ui.theme.onBackgroundColorFilter
import kotlinx.coroutines.ExperimentalCoroutinesApi
import kotlinx.coroutines.flow.combine
import kotlinx.coroutines.flow.flatMapLatest
import kotlinx.coroutines.flow.flowOf
@OptIn(ExperimentalCoroutinesApi::class)
@Composable
private fun rememberLocalBlossomBridgeForProfilePics(): Boolean {
val sessionManager =
try {
Amethyst.instance.sessionManager
} catch (e: UninitializedPropertyAccessException) {
return false
}
val probe = Amethyst.instance.localBlossomCacheProbe
val flow =
remember {
combine(
sessionManager.accountContent.flatMapLatest { state ->
if (state is AccountState.LoggedIn) state.account.settings.useLocalBlossomCache else flowOf(false)
},
probe.available,
) { toggle, probeUp -> toggle && probeUp }
}
val state by flow.collectAsStateWithLifecycle(initialValue = false)
return state
}
@Composable
fun RobohashAsyncImage(
@@ -120,21 +93,16 @@ fun RobohashFallbackAsyncImage(
loadRobohash: Boolean,
autoPlayGif: Boolean = true,
) {
val useBridge = rememberLocalBlossomBridgeForProfilePics()
val bridgedModel =
remember(model, robot, useBridge) {
bridgeProfilePictureUrl(model, useBridge, robot)
}
if (bridgedModel != null && loadProfilePicture && isAnimatedMediaUrl(bridgedModel)) {
if (model != null && loadProfilePicture && isAnimatedMediaUrl(model)) {
GifProfilePicture(
userHex = robot,
userPicture = bridgedModel,
userPicture = model,
contentDescription = contentDescription,
modifier = modifier,
loadRobohash = loadRobohash,
autoPlay = autoPlayGif,
)
} else if (bridgedModel != null && loadProfilePicture) {
} else if (model != null && loadProfilePicture) {
val fallbackPainter =
if (loadRobohash) {
rememberVectorPainter(
@@ -154,10 +122,10 @@ fun RobohashFallbackAsyncImage(
// file://) would fail there. Route only remote http(s) pictures through the thumbnail
// cache; hand local/content URIs to Coil's native fetchers, which load them directly.
model =
if (bridgedModel.startsWith("http://", ignoreCase = true) || bridgedModel.startsWith("https://", ignoreCase = true)) {
ProfilePictureUrl(bridgedModel)
if (model.startsWith("http://", ignoreCase = true) || model.startsWith("https://", ignoreCase = true)) {
ProfilePictureUrl(model)
} else {
bridgedModel
model
},
contentDescription = contentDescription,
modifier = modifier,
@@ -237,9 +205,25 @@ fun GifProfilePicture(
)
}
val context = LocalContext.current
// Animated avatars skip ProfilePictureFetcher (its thumbnail cache would flatten them), so
// they carry the profile-picture marker themselves for the local Blossom cache bridge.
val model =
remember(userPicture) {
ImageRequest
.Builder(context)
.data(userPicture)
.httpHeaders(
NetworkHeaders
.Builder()
.set(LocalBlossomCacheRedirectInterceptor.MEDIA_HEADER, LocalBlossomCacheRedirectInterceptor.PROFILE_PICTURE)
.build(),
).build()
}
Box(modifier = modifier) {
SubcomposeAsyncImage(
model = userPicture,
model = model,
contentDescription = contentDescription,
contentScale = ContentScale.Crop,
modifier = Modifier.fillMaxSize(),
@@ -27,6 +27,7 @@ import androidx.core.content.FileProvider
import com.vitorpamplona.amethyst.Amethyst
import com.vitorpamplona.amethyst.commons.richtext.mimeTypeMap
import com.vitorpamplona.amethyst.commons.richtext.normalizeMimeType
import com.vitorpamplona.amethyst.service.images.BlossomFetcher
import com.vitorpamplona.quartz.utils.Log
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.withContext
@@ -93,7 +94,7 @@ object ShareHelper {
): Pair<Uri, String> =
withContext(Dispatchers.IO) {
// Safely get snapshot and file
Amethyst.instance.diskCache.openSnapshot(imageUrl)?.use { snapshot ->
Amethyst.instance.diskCache.openSnapshot(BlossomFetcher.diskCacheKey(imageUrl))?.use { snapshot ->
val file = snapshot.data.toFile()
// Determine file extension and prepare sharable file
@@ -64,7 +64,6 @@ import androidx.compose.ui.util.lerp
import androidx.compose.ui.window.Dialog
import androidx.compose.ui.window.DialogProperties
import androidx.core.net.toUri
import androidx.lifecycle.compose.collectAsStateWithLifecycle
import com.vitorpamplona.amethyst.Amethyst
import com.vitorpamplona.amethyst.commons.resources.Res
import com.vitorpamplona.amethyst.commons.resources.failed_to_save_the_image
@@ -81,7 +80,6 @@ import com.vitorpamplona.amethyst.commons.richtext.MediaUrlContent
import com.vitorpamplona.amethyst.commons.richtext.MediaUrlImage
import com.vitorpamplona.amethyst.commons.richtext.MediaUrlPdf
import com.vitorpamplona.amethyst.commons.richtext.MediaUrlVideo
import com.vitorpamplona.amethyst.commons.richtext.toCoilModel
import com.vitorpamplona.amethyst.commons.ui.loadStringRes
import com.vitorpamplona.amethyst.model.MediaAspectRatioCache
import com.vitorpamplona.amethyst.service.playback.composable.VideoViewInner
@@ -505,12 +503,6 @@ private fun RenderImageOrVideo(
}
val ratio = content.dim?.aspectRatioOrNull() ?: MediaAspectRatioCache.get(content.url)
val useLocalBlossomBridge by accountViewModel.useLocalBlossomBridge.collectAsStateWithLifecycle()
val bridgedUrl =
remember(content.url, useLocalBlossomBridge) {
content.toCoilModel(useLocalBlossomBridge)
}
val modifier =
if (ratio != null) {
Modifier.aspectRatio(ratio)
@@ -520,7 +512,7 @@ private fun RenderImageOrVideo(
Box(modifier, contentAlignment = Alignment.Center) {
VideoViewInner(
videoUri = bridgedUrl,
videoUri = content.url,
mimeType = content.mimeType,
aspectRatio = ratio,
title = content.description,
@@ -66,7 +66,6 @@ import androidx.compose.ui.text.style.TextOverflow
import androidx.compose.ui.text.withStyle
import androidx.compose.ui.unit.dp
import androidx.core.net.toUri
import androidx.lifecycle.compose.collectAsStateWithLifecycle
import androidx.lifecycle.viewModelScope
import coil3.compose.AsyncImage
import coil3.compose.AsyncImagePainter
@@ -105,11 +104,11 @@ import com.vitorpamplona.amethyst.commons.richtext.MediaUrlImage
import com.vitorpamplona.amethyst.commons.richtext.MediaUrlPdf
import com.vitorpamplona.amethyst.commons.richtext.MediaUrlVideo
import com.vitorpamplona.amethyst.commons.richtext.RichTextParser
import com.vitorpamplona.amethyst.commons.richtext.toCoilModel
import com.vitorpamplona.amethyst.commons.service.image.placeholderModel
import com.vitorpamplona.amethyst.commons.ui.components.LoadingAnimation
import com.vitorpamplona.amethyst.commons.ui.loadStringRes
import com.vitorpamplona.amethyst.model.MediaAspectRatioCache
import com.vitorpamplona.amethyst.service.images.BlossomFetcher
import com.vitorpamplona.amethyst.service.playback.composable.VideoView
import com.vitorpamplona.amethyst.service.uploads.blossom.bud10.openBlossomUriAsIntent
import com.vitorpamplona.amethyst.ui.actions.CrossfadeIfEnabled
@@ -197,26 +196,20 @@ fun ZoomableContentView(
sourceBounds = coordinates.boundsInWindow()
}
val useLocalBlossomBridge by accountViewModel.useLocalBlossomBridge.collectAsStateWithLifecycle()
when (content) {
is MediaUrlImage -> {
val ratio = content.dim?.aspectRatioOrNull() ?: MediaAspectRatioCache.get(content.url)
val bridgedUrl =
remember(content.url, useLocalBlossomBridge) {
content.toCoilModel(useLocalBlossomBridge)
}
ContentWarningGate(
isSensitive = content.contentWarning != null,
reasons = setOfNotNull(content.contentWarning),
preloadUrls = listOf(bridgedUrl),
preloadUrls = listOf(content.url),
accountViewModel = accountViewModel,
modifier = mediaSizingModifier(ratio, contentScale),
backdrop = (content.thumbhash ?: content.blurhash)?.let { { BlurhashBackdrop(content.blurhash, content.description, content.thumbhash) } },
) {
if (content.isAnimatedMedia()) {
GifVideoView(
videoUri = bridgedUrl,
videoUri = content.url,
contentDescription = content.description,
dimensions = content.dim,
blurhash = content.blurhash,
@@ -251,10 +244,6 @@ fun ZoomableContentView(
content.dim?.aspectRatioOrNull()
?: MediaAspectRatioCache.get(content.url)
?: fallbackRatio
val bridgedUrl =
remember(content.url, useLocalBlossomBridge) {
content.toCoilModel(useLocalBlossomBridge)
}
ContentWarningGate(
isSensitive = content.contentWarning != null,
reasons = setOfNotNull(content.contentWarning),
@@ -274,7 +263,7 @@ fun ZoomableContentView(
contentAlignment = Alignment.Center,
) {
VideoView(
videoUri = bridgedUrl,
videoUri = content.url,
mimeType = content.mimeType,
title = content.description,
artworkUri = content.artworkUri,
@@ -541,22 +530,17 @@ fun UrlImageView(
}
val context = LocalContext.current
val useLocalBlossomBridge by accountViewModel.useLocalBlossomBridge.collectAsStateWithLifecycle()
val bridgedUrl =
remember(content.url, useLocalBlossomBridge) {
content.toCoilModel(useLocalBlossomBridge)
}
val imageModel =
if (fullResolution) {
remember(bridgedUrl, context) {
remember(content.url, context) {
ImageRequest
.Builder(context)
.data(bridgedUrl)
.data(content.url)
.size(Size.ORIGINAL)
.build()
}
} else {
bridgedUrl
content.url
}
CrossfadeIfEnabled(targetState = showImage.value, contentAlignment = Alignment.Center, accountViewModel = accountViewModel) {
@@ -1278,15 +1262,9 @@ private suspend fun shareLocalVideoFile(
private fun verifyHash(content: MediaUrlContent): Boolean? {
if (content.hash == null) return null
val keys = mutableListOf(content.url)
val bridged = content.toCoilModel(true)
if (bridged != content.url) keys.add(bridged)
for (key in keys) {
Amethyst.instance.diskCache.openSnapshot(key)?.use { snapshot ->
val (hashBytes, _) = sha256StreamWithCount(snapshot.data.toFile().inputStream())
return hashBytes.toHexKey() == content.hash
}
Amethyst.instance.diskCache.openSnapshot(BlossomFetcher.diskCacheKey(content.url))?.use { snapshot ->
val (hashBytes, _) = sha256StreamWithCount(snapshot.data.toFile().inputStream())
return hashBytes.toHexKey() == content.hash
}
return null
@@ -276,32 +276,10 @@ class AccountViewModel(
val feedStates = AccountFeedContentStates(account, viewModelScope)
/**
* `true` when feed/note media (images and videos in `MediaUrlContent`)
* should be routed through the local Blossom cache. Requires the master
* toggle on, the probe up, AND the profile-pictures-only restriction
* to be off.
* `true` when the local Blossom cache is enabled and the probe sees it up. Only drives the
* "detected" chip in settings: the routing itself happens in LocalBlossomCacheRedirectInterceptor.
*/
val useLocalBlossomBridge: StateFlow<Boolean> =
try {
combine(
account.settings.useLocalBlossomCache,
account.settings.localBlossomCacheProfilePicturesOnly,
Amethyst.instance.localBlossomCacheProbe.available,
) { toggle, profileOnly, probeUp -> toggle && probeUp && !profileOnly }.stateIn(
viewModelScope,
SharingStarted.Eagerly,
false,
)
} catch (e: UninitializedPropertyAccessException) {
MutableStateFlow(false)
}
/**
* `true` when profile pictures should be routed through the local
* Blossom cache. Requires only the master toggle and the probe to be
* up; the profile-pictures-only restriction does not gate this flow.
*/
val useLocalBlossomBridgeForProfilePics: StateFlow<Boolean> =
val localBlossomCacheDetected: StateFlow<Boolean> =
try {
combine(
account.settings.useLocalBlossomCache,
@@ -35,7 +35,6 @@ import androidx.compose.ui.Modifier
import androidx.compose.ui.draw.clip
import androidx.compose.ui.graphics.Color
import androidx.compose.ui.layout.ContentScale
import androidx.lifecycle.compose.collectAsStateWithLifecycle
import androidx.media3.common.util.UnstableApi
import coil3.compose.AsyncImagePainter
import coil3.compose.SubcomposeAsyncImage
@@ -50,7 +49,6 @@ import com.vitorpamplona.amethyst.commons.richtext.MediaUrlImage
import com.vitorpamplona.amethyst.commons.richtext.MediaUrlVideo
import com.vitorpamplona.amethyst.commons.richtext.RichTextParser.Companion.isHlsMimeType
import com.vitorpamplona.amethyst.commons.richtext.RichTextParser.Companion.isVideoUrl
import com.vitorpamplona.amethyst.commons.richtext.toCoilModel
import com.vitorpamplona.amethyst.commons.ui.components.LoadingAnimation
import com.vitorpamplona.amethyst.service.playback.composable.mediaitem.isHlsMedia
import com.vitorpamplona.amethyst.service.relayClient.reqCommand.event.observeNote
@@ -256,16 +254,11 @@ fun UrlImageView(
val isVideo = content is MediaUrlVideo
val artworkUri = (content as? MediaUrlVideo)?.artworkUri
val useLocalBlossomBridge by accountViewModel.useLocalBlossomBridge.collectAsStateWithLifecycle()
// Coil's VideoFrameDecoder can extract a frame from .mp4/.webm but not from an HLS .m3u8
// playlist (it's a text manifest). For an HLS video without a separate artwork URL, sending
// the playlist to SubcomposeAsyncImage just produces an Error state and a stand-in icon.
// Skip the fetch in that case and render blurhash + play overlay directly.
val bridgedUrl =
remember(content.url, useLocalBlossomBridge) {
content.toCoilModel(useLocalBlossomBridge)
}
val imageModelUrl = artworkUri ?: bridgedUrl
val imageModelUrl = artworkUri ?: content.url
val canLoadAsImage = !isVideo || artworkUri != null || !isHlsMedia(content.url, content.mimeType)
CrossfadeIfEnabled(targetState = showImage.value, contentAlignment = Alignment.Center, accountViewModel = accountViewModel) {
@@ -206,7 +206,13 @@ class LocalBlossomCacheRedirectInterceptorTest {
url: String,
captured: MutableList<String>,
): Interceptor.Chain {
val request = Request.Builder().url(url.toHttpUrl()).build()
// Marked the way the media loaders mark their downloads; only those are bridged.
val request =
Request
.Builder()
.url(url.toHttpUrl())
.header(LocalBlossomCacheRedirectInterceptor.MEDIA_HEADER, LocalBlossomCacheRedirectInterceptor.MEDIA)
.build()
return Proxy.newProxyInstance(
Interceptor.Chain::class.java.classLoader,
arrayOf(Interceptor.Chain::class.java),
@@ -48,9 +48,16 @@ import com.vitorpamplona.quartz.nip57Zaps.LnZapRequestEvent
import io.mockk.every
import io.mockk.mockk
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.ExperimentalCoroutinesApi
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.test.runTest
import okhttp3.Interceptor
import okhttp3.Protocol
import okhttp3.Request
import okhttp3.Response
import okhttp3.ResponseBody.Companion.toResponseBody
import org.junit.Assert.assertEquals
import org.junit.Assert.assertFalse
import org.junit.Assert.assertNull
@@ -387,6 +394,48 @@ class RadioBurstEstimatorTest {
}
class LoopbackExclusionTest {
@get:Rule
val temp = TemporaryFolder()
/** Runs [UsageCountingInterceptor] over a chain whose downstream sent [sentUrl] for [url]. */
private fun countedRequests(
url: String,
sentUrl: String,
): Long? {
val accountant = ResourceUsageAccountant(ResourceUsageStore(File(temp.root, "u.json")), CoroutineScope(Dispatchers.Unconfined), epochDay = { 1L })
val interceptor = UsageCountingInterceptor(accountant, isMobile = { true }, isForeground = { true }, nowMs = { 0L })
val request = Request.Builder().url(url).build()
val chain = mockk<Interceptor.Chain>()
every { chain.request() } returns request
every { chain.proceed(any()) } answers {
Response
.Builder()
.request(request.newBuilder().url(sentUrl).build())
.protocol(Protocol.HTTP_1_1)
.code(200)
.message("OK")
.body("".toResponseBody(null))
.build()
}
interceptor.intercept(chain).close()
return runBlocking { accountant.allDaysIncludingLive() }[1L].orEmpty()[UsageKeys.netReqs(UsageKeys.ROLE_OTHER, mobile = true, foreground = true)]
}
@Test
fun requestRewrittenToTheLocalCacheIsNotCounted() {
// This interceptor is outermost: it sees the CDN URL, but the request that actually went
// out (LocalBlossomCacheRedirectInterceptor's rewrite) never touched the radio.
val sha = "b1674191a88ec5cdd733e4240a81803105dc412d6c6708d53ab94fc248f4f553"
assertNull(countedRequests("https://cdn.example.com/$sha.jpg", "http://127.0.0.1:24242/$sha.jpg"))
}
@Test
fun requestToARealHostIsCounted() {
assertEquals(1L, countedRequests("https://cdn.example.com/a.jpg", "https://cdn.example.com/a.jpg"))
}
@Test
fun loopbackHostsAreRecognizedAndRealHostsAreNot() {
assertTrue(UsageCountingInterceptor.isLoopback("127.0.0.1"))
@@ -0,0 +1,95 @@
/*
* Copyright (c) 2025 Vitor Pamplona
*
* Permission is hereby granted, free of charge, to any person obtaining a copy of
* this software and associated documentation files (the "Software"), to deal in
* the Software without restriction, including without limitation the rights to use,
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
* Software, and to permit persons to whom the Software is furnished to do so,
* subject to the following conditions:
*
* The above copyright notice and this permission notice shall be included in all
* copies or substantial portions of the Software.
*
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
*/
package com.vitorpamplona.amethyst.service.uploads.blossom.bud10
import com.vitorpamplona.amethyst.commons.service.http.IRoleBasedHttpClientBuilder
import com.vitorpamplona.quartz.nipB7Blossom.BlossomServersEvent
import io.mockk.mockk
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.awaitCancellation
import kotlinx.coroutines.cancel
import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.flow
import kotlinx.coroutines.launch
import kotlinx.coroutines.runBlocking
import org.junit.After
import org.junit.Assert.assertEquals
import org.junit.Assert.assertNull
import org.junit.Test
import java.util.concurrent.atomic.AtomicInteger
class BlossomServerResolverTest {
private val sha = "b1674191a88ec5cdd733e4240a81803105dc412d6c6708d53ab94fc248f4f553"
private val author = "460c25e682fda7832b52d1f22d3d22b3176d972f60dcdc3212ed8c92ef85065c"
private val scope = CoroutineScope(SupervisorJob() + Dispatchers.IO)
@After
fun tearDown() = scope.cancel()
/** A resolver whose author server-list lookups are counted and answered by [serverList]. */
private fun resolver(
lookups: AtomicInteger,
serverList: () -> Flow<BlossomServersEvent>,
) = BlossomServerResolver(
loggedInUsers = { emptyList() },
blossomServers = { addresses ->
lookups.incrementAndGet()
addresses.map { serverList() }
},
httpClientBuilder = mockk<IRoleBasedHttpClientBuilder>(relaxed = true),
scope = scope,
)
@Test
fun concurrentLookupsOfTheSameUriShareOneResolution() =
runBlocking {
val lookups = AtomicInteger(0)
// The author's server list never arrives, like a relay that is slow to answer.
val resolver = resolver(lookups) { flow { awaitCancellation() } }
val uri = "blossom:$sha.jpg?as=$author"
val callers = List(5) { scope.launch { resolver.findServers(uri) } }
delay(300)
assertEquals(1, lookups.get())
callers.forEach { it.cancel() }
}
@Test
fun anUnresolvableUriIsNotRetriedUntilTheMissExpiresOrCachesAreCleared() =
runBlocking {
val lookups = AtomicInteger(0)
// No hints and no authors: resolves to nothing right away.
val resolver = resolver(lookups) { flow { } }
val uri = "blossom:$sha.jpg"
assertNull(resolver.findServers(uri))
assertNull(resolver.findServers(uri))
assertEquals(1, lookups.get())
resolver.clearCaches()
assertNull(resolver.findServers(uri))
assertEquals(2, lookups.get())
}
}
@@ -1,225 +0,0 @@
/*
* Copyright (c) 2025 Vitor Pamplona
*
* Permission is hereby granted, free of charge, to any person obtaining a copy of
* this software and associated documentation files (the "Software"), to deal in
* the Software without restriction, including without limitation the rights to use,
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
* Software, and to permit persons to whom the Software is furnished to do so,
* subject to the following conditions:
*
* The above copyright notice and this permission notice shall be included in all
* copies or substantial portions of the Software.
*
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
*/
package com.vitorpamplona.amethyst.commons.richtext
import com.vitorpamplona.quartz.nipB7Blossom.BlossomUri
private val sha256HexRegex = Regex("[0-9a-f]{64}")
private val blossomLastSegmentRegex = Regex("^([0-9a-fA-F]{64})(?:\\.[^./]+)?$")
/**
* Converts this media content into a Coil/ExoPlayer-friendly model string.
*
* When the local-Blossom-cache bridge is active and the content has a
* known sha256 hash, returns a `blossom:<sha256>.<ext>?xs=<originalHostBase>&as=<authorPubKey>`
* URI. The Coil pipeline recognises the scheme and routes the request
* through `BlossomServerResolver`, which short-circuits to the local cache
* at `127.0.0.1:24242`.
*
* Otherwise (bridge off, no hash, hash invalid, already a `blossom:` URI,
* or a live stream) returns the original URL unchanged so today's
* direct-to-CDN behaviour is preserved.
*/
fun MediaUrlContent.toCoilModel(useLocalBlossomBridge: Boolean): String =
bridgeUrl(
url = url,
useBridge = useLocalBlossomBridge,
mimeType = mimeType,
authorPubKey = authorPubKey,
skipBridge = this is MediaUrlVideo && isLiveStream,
)
/**
* Bridge entry point for raw URL strings (e.g. profile pictures) that go
* through a Coil fetcher routed by type (`ProfilePictureUrl`) and therefore
* bypass [com.vitorpamplona.quartz.nipB7Blossom.BlossomUri] processing
* entirely.
*
* Returns a direct `http://127.0.0.1:24242/<sha>.<ext>?xs=<host>&as=<pubkey>`
* URL so the request can flow through `NetworkFetcher` unchanged. Falls
* back to the original URL when no sha256 can be recovered from the path
* or when the bridge is off.
*
* @param localCacheBase the local Blossom cache origin (default `http://127.0.0.1:24242`).
* @param authorPubKey 64-char lowercase hex pubkey appended as `as=` so the
* cache can consult that author's BUD-03 server list on miss.
*/
fun bridgeProfilePictureUrl(
url: String?,
useBridge: Boolean,
authorPubKey: String? = null,
localCacheBase: String = DEFAULT_LOCAL_CACHE_BASE,
): String? {
if (url == null) return null
if (!useBridge) return url
if (url.startsWith("blossom:", ignoreCase = true)) return url
if (!url.startsWith("http://", ignoreCase = true) && !url.startsWith("https://", ignoreCase = true)) return url
val sha = extractSha256FromUrlPath(url) ?: return url
val ext = guessExtension(url, null)
val serverBase = extractServerBase(url, sha) ?: return url
val params =
buildList {
add("xs=${percentEncode(serverBase)}")
authorPubKey
?.lowercase()
?.takeIf { sha256HexRegex.matches(it) }
?.let { add("as=$it") }
}
return "${localCacheBase.removeSuffix("/")}/$sha.$ext?${params.joinToString("&")}"
}
const val DEFAULT_LOCAL_CACHE_BASE = "http://127.0.0.1:24242"
private fun bridgeUrl(
url: String,
useBridge: Boolean,
mimeType: String?,
authorPubKey: String?,
skipBridge: Boolean,
): String {
if (!useBridge || skipBridge) return url
if (url.startsWith("blossom:", ignoreCase = true)) return url
if (!url.startsWith("http://", ignoreCase = true) && !url.startsWith("https://", ignoreCase = true)) return url
// The local Blossom cache fetches `<xs>/<sha>.<ext>` on miss per BUD-01,
// which only works when the upstream URL is itself BUD-01 layout — the
// file at `<xs>/<sha>.<ext>` is the one named in the URL path, not the
// imeta `x` hash (which on resizing CDNs may identify a different blob
// than the URL filename, e.g. `x` = post-resize, `ox` = original).
// Always use the URL's sha; never trust the imeta override.
val sha = extractSha256FromUrlPath(url) ?: return url
val ext = guessExtension(url, mimeType)
val serverBase = extractServerBase(url, sha) ?: return url
val authors =
authorPubKey
?.lowercase()
?.takeIf { sha256HexRegex.matches(it) }
?.let { listOf(it) }
?: emptyList()
return BlossomUri(
sha256 = sha,
extension = ext,
servers = listOf(serverBase),
authors = authors,
size = null,
).toUriString()
}
private fun percentEncode(input: String): String {
val sb = StringBuilder(input.length)
for (c in input) {
if (c.isLetterOrDigit() || c in "-._~") {
sb.append(c)
} else {
for (b in c.toString().encodeToByteArray()) {
sb.append('%')
sb.append((b.toInt() and 0xFF).toString(16).padStart(2, '0').uppercase())
}
}
}
return sb.toString()
}
private fun extractSha256FromUrlPath(url: String): String? {
// Per Blossom (BUD-01) the last path segment must be exactly
// `<sha256>` or `<sha256>.<ext>`. URLs whose filename merely embeds a
// 64-char hex (e.g. "nostr.build_<sha>.jpg") aren't Blossom blobs and
// the bridge must leave them alone — rewriting them would point the
// local cache at a fallback `xs=` server that doesn't host the blob.
val pathPart = url.substringBefore('?').substringBefore('#')
val lastSegment = pathPart.substringAfterLast('/')
val match = blossomLastSegmentRegex.matchEntire(lastSegment) ?: return null
return match.groupValues[1].lowercase()
}
private fun guessExtension(
url: String,
mimeType: String?,
): String {
val pathPart = url.substringBefore('?').substringBefore('#')
val lastDot = pathPart.lastIndexOf('.')
val lastSlash = pathPart.lastIndexOf('/')
if (lastDot > lastSlash && lastDot >= 0) {
val ext = pathPart.substring(lastDot + 1).lowercase()
if (ext.isNotEmpty() && ext.length <= 8 && ext.all { it.isLetterOrDigit() }) {
return ext
}
}
if (mimeType != null) {
for ((extension, mt) in mimeTypeMap) {
if (mt.equals(mimeType, ignoreCase = true)) return extension
}
}
return "bin"
}
/**
* Returns the URL prefix that the local Blossom cache should append `/<sha>`
* to in order to reach the original blob, preserving any path prefix the
* upstream CDN uses (e.g. `https://cdn.nostr.build/i` for nostr.build's
* `/i/<sha>` scheme). Falls back to scheme+host when the sha can't be
* located in the path.
*/
private fun extractServerBase(
url: String,
sha: String,
): String? {
val pathPart = url.substringBefore('?').substringBefore('#')
val schemeEnd = pathPart.indexOf("://")
if (schemeEnd < 0) return null
val hostStart = schemeEnd + 3
if (hostStart >= pathPart.length) return null
val shaIndex = pathPart.lastIndexOf(sha, ignoreCase = true)
if (shaIndex >= 0) {
// Anchor on the slash immediately preceding the sha so the cache
// can append "/<sha>" verbatim per the local-blossom-cache spec.
val slashBeforeSha = pathPart.lastIndexOf('/', shaIndex - 1)
if (slashBeforeSha > hostStart - 1) {
return pathPart.substring(0, slashBeforeSha)
}
}
return extractHostBase(pathPart)
}
private fun extractHostBase(url: String): String? {
val schemeEnd = url.indexOf("://")
if (schemeEnd < 0) return null
val afterScheme = schemeEnd + 3
var end = url.length
for (i in afterScheme until url.length) {
val c = url[i]
if (c == '/' || c == '?' || c == '#') {
end = i
break
}
}
if (end <= afterScheme) return null
return url.substring(0, end)
}
@@ -1,273 +0,0 @@
/*
* Copyright (c) 2025 Vitor Pamplona
*
* Permission is hereby granted, free of charge, to any person obtaining a copy of
* this software and associated documentation files (the "Software"), to deal in
* the Software without restriction, including without limitation the rights to use,
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
* Software, and to permit persons to whom the Software is furnished to do so,
* subject to the following conditions:
*
* The above copyright notice and this permission notice shall be included in all
* copies or substantial portions of the Software.
*
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
*/
package com.vitorpamplona.amethyst.commons.richtext
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertTrue
class MediaUrlContentExtTest {
private val sha = "b1674191a88ec5cdd733e4240a81803105dc412d6c6708d53ab94fc248f4f553"
@Test
fun bridgeOffReturnsOriginalUrl() {
val image = MediaUrlImage(url = "https://cdn.example.com/$sha.jpg", hash = sha)
assertEquals("https://cdn.example.com/$sha.jpg", image.toCoilModel(useLocalBlossomBridge = false))
}
@Test
fun nullHashReturnsOriginalUrl() {
val image = MediaUrlImage(url = "https://cdn.example.com/foo.jpg", hash = null)
assertEquals("https://cdn.example.com/foo.jpg", image.toCoilModel(useLocalBlossomBridge = true))
}
@Test
fun invalidHashReturnsOriginalUrl() {
val image = MediaUrlImage(url = "https://cdn.example.com/foo.jpg", hash = "not-hex")
assertEquals("https://cdn.example.com/foo.jpg", image.toCoilModel(useLocalBlossomBridge = true))
}
@Test
fun blossomUriReturnedUnchanged() {
val image = MediaUrlImage(url = "blossom:$sha.jpg?xs=https://cdn.example.com", hash = sha)
assertEquals("blossom:$sha.jpg?xs=https://cdn.example.com", image.toCoilModel(useLocalBlossomBridge = true))
}
@Test
fun liveStreamReturnsOriginalUrl() {
val video = MediaUrlVideo(url = "https://stream.example.com/play.m3u8", hash = sha, isLiveStream = true)
assertEquals("https://stream.example.com/play.m3u8", video.toCoilModel(useLocalBlossomBridge = true))
}
@Test
fun bridgeOnRewritesPlainHttpsUrl() {
val image = MediaUrlImage(url = "https://nostr.build/i/abc/$sha.jpg", hash = sha)
val result = image.toCoilModel(useLocalBlossomBridge = true)
assertEquals("blossom:$sha.jpg?xs=https://nostr.build/i/abc", result)
}
@Test
fun bridgeOnPreservesNostrBuildPathPrefix() {
val image = MediaUrlImage(url = "https://cdn.nostr.build/i/$sha.jpg", hash = sha)
val result = image.toCoilModel(useLocalBlossomBridge = true)
assertEquals("blossom:$sha.jpg?xs=https://cdn.nostr.build/i", result)
}
@Test
fun bridgeOnFlatBlossomPathYieldsHostOnlyXs() {
val image = MediaUrlImage(url = "https://blossom.primal.net/$sha.jpg", hash = sha)
val result = image.toCoilModel(useLocalBlossomBridge = true)
assertEquals("blossom:$sha.jpg?xs=https://blossom.primal.net", result)
}
@Test
fun bridgeOnInfersExtensionFromMimeType() {
// BUD-01 allows `<sha>` without an extension; mimeType supplies one.
val image = MediaUrlImage(url = "https://nostr.build/i/$sha", hash = sha, mimeType = "image/png")
val result = image.toCoilModel(useLocalBlossomBridge = true)
assertTrue(result.startsWith("blossom:$sha.png?xs="), "expected png extension from mime, got $result")
}
@Test
fun bridgeOnFallsBackToBinExtension() {
val image = MediaUrlImage(url = "https://nostr.build/i/$sha", hash = sha)
val result = image.toCoilModel(useLocalBlossomBridge = true)
assertTrue(result.startsWith("blossom:$sha.bin?xs="), "expected bin extension fallback, got $result")
}
@Test
fun nonHttpUrlReturnsOriginal() {
val image = MediaUrlImage(url = "ftp://example.com/file", hash = sha)
assertEquals("ftp://example.com/file", image.toCoilModel(useLocalBlossomBridge = true))
}
@Test
fun uppercaseHashNormalisedToLowercase() {
val image = MediaUrlImage(url = "https://cdn.example.com/${sha.uppercase()}.jpg", hash = sha.uppercase())
val result = image.toCoilModel(useLocalBlossomBridge = true)
assertTrue(result.startsWith("blossom:$sha.jpg?xs="), "expected lowercase sha, got $result")
}
@Test
fun authorPubKeyAddedAsAsParam() {
val authorPub = "a8f3721a0dc1b4d5c12f4cc7c54ae14071eb9c1b4f9b2cf0d4ab22c0e9f0c7e5"
val image =
MediaUrlImage(
url = "https://cdn.example.com/$sha.jpg",
hash = sha,
authorPubKey = authorPub,
)
val result = image.toCoilModel(useLocalBlossomBridge = true)
assertEquals("blossom:$sha.jpg?xs=https://cdn.example.com&as=$authorPub", result)
}
@Test
fun invalidAuthorPubKeyDropped() {
val image =
MediaUrlImage(
url = "https://cdn.example.com/$sha.jpg",
hash = sha,
authorPubKey = "not-a-pubkey",
)
val result = image.toCoilModel(useLocalBlossomBridge = true)
assertEquals("blossom:$sha.jpg?xs=https://cdn.example.com", result)
}
@Test
fun bridgeOnSkipsNonBud01UrlEvenWithImetaHash() {
// The imeta `x` hash refers to a blob whose canonical Blossom location
// is /<sha>.<ext>, but the upstream URL serves it under a different
// path (https://i.nostr.build/M5AwJ.gif). Routing this through the
// local cache would set xs=https://i.nostr.build, and the cache would
// fetch https://i.nostr.build/<sha>.gif on miss, which 404s.
val url = "https://i.nostr.build/M5AwJ.gif"
val image = MediaUrlImage(url = url, hash = sha)
assertEquals(url, image.toCoilModel(useLocalBlossomBridge = true))
}
@Test
fun bridgeProfilePictureUrlSkipsNonBud01UrlEvenWithImetaHash() {
val url = "https://i.nostr.build/M5AwJ.gif"
assertEquals(url, bridgeProfilePictureUrl(url, useBridge = true))
}
@Test
fun bridgeProfilePictureUrlNullReturnsNull() {
assertEquals(null, bridgeProfilePictureUrl(null, useBridge = true))
}
@Test
fun bridgeProfilePictureUrlOffReturnsOriginal() {
assertEquals(
"https://cdn.example.com/avatar.jpg",
bridgeProfilePictureUrl("https://cdn.example.com/avatar.jpg", useBridge = false),
)
}
@Test
fun bridgeProfilePictureUrlExtractsShaFromPath() {
val url = "https://nostr.build/i/$sha.jpg"
val authorPub = "a8f3721a0dc1b4d5c12f4cc7c54ae14071eb9c1b4f9b2cf0d4ab22c0e9f0c7e5"
assertEquals(
"http://127.0.0.1:24242/$sha.jpg?xs=https%3A%2F%2Fnostr.build%2Fi&as=$authorPub",
bridgeProfilePictureUrl(url, useBridge = true, authorPubKey = authorPub),
)
}
@Test
fun bridgeProfilePictureUrlPreservesNostrBuildPath() {
val url = "https://cdn.nostr.build/i/$sha.jpg"
assertEquals(
"http://127.0.0.1:24242/$sha.jpg?xs=https%3A%2F%2Fcdn.nostr.build%2Fi",
bridgeProfilePictureUrl(url, useBridge = true),
)
}
@Test
fun bridgeProfilePictureUrlNoShaInPathReturnsOriginal() {
val url = "https://nostr.build/avatar.jpg"
assertEquals(url, bridgeProfilePictureUrl(url, useBridge = true))
}
@Test
fun bridgeProfilePictureUrlBlossomUriReturnedUnchanged() {
val uri = "blossom:$sha.jpg?xs=https://nostr.build"
assertEquals(uri, bridgeProfilePictureUrl(uri, useBridge = true))
}
@Test
fun bridgeOnRewritesShaInLastPathSegmentWithHexPrefix() {
// share.yabu.me layout: <cache-prefix-sha>/<blob-sha>.<ext>
val image =
MediaUrlImage(
url = "https://share.yabu.me/84b0c46ab699ac35eb2ca286470b85e081db2087cdef63932236c397417782f5/28fa4d999af6ae3e4e11bfc2727130ef1b3a13cc0f981e5a93c3996cb2f524e5.webp",
hash = null,
)
assertEquals(
"blossom:28fa4d999af6ae3e4e11bfc2727130ef1b3a13cc0f981e5a93c3996cb2f524e5.webp?xs=https://share.yabu.me/84b0c46ab699ac35eb2ca286470b85e081db2087cdef63932236c397417782f5",
image.toCoilModel(useLocalBlossomBridge = true),
)
}
@Test
fun bridgeProfilePictureUrlRewritesShaInLastPathSegmentWithHexPrefix() {
assertEquals(
"http://127.0.0.1:24242/28fa4d999af6ae3e4e11bfc2727130ef1b3a13cc0f981e5a93c3996cb2f524e5.webp?xs=https%3A%2F%2Fshare.yabu.me%2F84b0c46ab699ac35eb2ca286470b85e081db2087cdef63932236c397417782f5",
bridgeProfilePictureUrl(
"https://share.yabu.me/84b0c46ab699ac35eb2ca286470b85e081db2087cdef63932236c397417782f5/28fa4d999af6ae3e4e11bfc2727130ef1b3a13cc0f981e5a93c3996cb2f524e5.webp",
useBridge = true,
),
)
}
@Test
fun bridgeOnSkipsWhenLastSegmentIsNotSha() {
// Per BUD-01 the last segment is the blob; if it isn't a sha256, the
// URL isn't a Blossom blob even if an earlier segment is hex.
val url = "https://example.com/$sha/avatar.jpg"
val image = MediaUrlImage(url = url, hash = null)
assertEquals(url, image.toCoilModel(useLocalBlossomBridge = true))
}
@Test
fun bridgeProfilePictureUrlSkipsWhenLastSegmentIsNotSha() {
val url = "https://example.com/$sha/avatar.jpg"
assertEquals(url, bridgeProfilePictureUrl(url, useBridge = true))
}
@Test
fun bridgeOnSkipsWhenLastSegmentHasNonHexPrefixBeforeSha() {
// nostr.build /i/ layout: <prefix>_<sha>.<ext>. The hex inside the
// filename isn't a Blossom blob per BUD-01 — the last segment must
// be exactly <sha> or <sha>.<ext>.
val url = "https://nostr.build/i/nostr.build_$sha.jpg"
val image = MediaUrlImage(url = url, hash = null)
assertEquals(url, image.toCoilModel(useLocalBlossomBridge = true))
}
@Test
fun bridgeProfilePictureUrlSkipsWhenLastSegmentHasNonHexPrefixBeforeSha() {
val url = "https://nostr.build/i/nostr.build_$sha.jpg"
assertEquals(url, bridgeProfilePictureUrl(url, useBridge = true))
}
@Test
fun bridgeOnSkipsWhenLastSegmentHasSuffixAfterSha() {
val url = "https://example.com/${sha}_thumb.jpg"
val image = MediaUrlImage(url = url, hash = null)
assertEquals(url, image.toCoilModel(useLocalBlossomBridge = true))
}
@Test
fun bridgeOnUsesUrlShaNotImetaWhenTheyDiffer() {
// On resizing CDNs the imeta `x` (post-resize) can differ from the
// `ox` (original) embedded in the URL. The upstream file is named
// after the URL's sha, so the cache request must use that — using
// the imeta `x` would point xs= at a non-existent path on miss.
val urlSha = "f24026b7281e598973a775adefb1b9a13b9f037a94ac98dd48ccc91b83f4b7b3"
val imetaX = "6932a918de1bfae3bf6611794ff54dd677013d22b760a9212117a0bd9079badf"
val image = MediaUrlImage(url = "https://image.nostr.build/$urlSha.png", hash = imetaX)
assertEquals(
"blossom:$urlSha.png?xs=https://image.nostr.build",
image.toCoilModel(useLocalBlossomBridge = true),
)
}
}
@@ -20,6 +20,7 @@
*/
package com.vitorpamplona.amethyst.commons.service.http
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.SurgeDns
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.flow.SharingStarted
@@ -41,7 +42,7 @@ class DualHttpClientManager(
keyCache: EncryptionKeyCache,
scope: CoroutineScope,
dns: SurgeDns,
shouldBridgeBlossomCache: (() -> Boolean)? = null,
shouldBridgeBlossomCache: ((profilePicture: Boolean) -> Boolean)? = null,
// Required (not nullable): every general-purpose HTTP client we mint must be
// wired into the OnionLocationCache so the app's `.onion`-routing behavior
// is uniform across image, upload, NIP-05, money, preview, and push roles.
@@ -53,8 +54,10 @@ class DualHttpClientManager(
// Signs BUD-01 read-auth to retry auth-gated Blossom downloads on 401.
// See [BlossomReadAuthInterceptor].
blossomReadAuth: Interceptor? = null,
// Told when the local Blossom cache refuses a connection. See [LocalBlossomCacheRedirectInterceptor].
onLocalBlossomCacheUnreachable: () -> Unit = {},
) : IHttpClientManager {
val factory = OkHttpClientFactory(keyCache, userAgent, dns, shouldBridgeBlossomCache, onionCache, usageInterceptor, blossomReadAuth)
val factory = OkHttpClientFactory(keyCache, userAgent, dns, shouldBridgeBlossomCache, onionCache, usageInterceptor, blossomReadAuth, onLocalBlossomCacheUnreachable)
val defaultHttpClient: StateFlow<OkHttpClient> =
combine(proxyPortProvider, isMobileDataProvider) { proxy, mobile ->
@@ -96,10 +99,23 @@ class DualHttpClientManager(
/**
* the okhttp can change on the manager without affecting other systems.
*
* [useProxy] is fixed when the factory is built (e.g. per video player pool), from the URL the
* caller knew at the time. The URL that is finally requested can differ: a `blossom:` URI is
* resolved to the local Blossom cache on `127.0.0.1:24242` only when the data source opens. Tor
* refuses to connect to loopback/private addresses, so those requests must never take the proxied
* client — the same rule RoleBasedHttpClientBuilder applies to URLs it sees upfront.
*/
class DynamicCallFactory(
val useProxy: Boolean,
val manager: DualHttpClientManager,
) : Call.Factory {
override fun newCall(request: Request): Call = manager.getHttpClient(useProxy).newCall(request)
override fun newCall(request: Request): Call = manager.getHttpClient(shouldUseProxy(useProxy, request.url.toString())).newCall(request)
companion object {
fun shouldUseProxy(
useProxy: Boolean,
url: String,
): Boolean = useProxy && !RelayUrlNormalizer.isLocalHost(url) && !RelayUrlNormalizer.isOverlayNetwork(url)
}
}
@@ -20,43 +20,107 @@
*/
package com.vitorpamplona.amethyst.commons.service.http
import okhttp3.Call
import okhttp3.HttpUrl
import okhttp3.HttpUrl.Companion.toHttpUrl
import okhttp3.Interceptor
import okhttp3.Request
import okhttp3.Response
import java.net.ConnectException
/**
* App-wide OkHttp interceptor that transparently rewrites HTTP requests for
* OkHttp interceptor that transparently rewrites media downloads of
* sha256-keyed blobs to a local Blossom cache running on `127.0.0.1:24242`,
* per https://github.com/hzrd149/blossom/blob/master/implementations/local-blossom-cache.md
*
* Activates when [shouldBridge] returns `true` AND the request URL contains
* a 64-char hex sha256 segment in its path AND the host isn't already
* `127.0.0.1`/`localhost`. The original scheme+host is appended as a `xs=`
* proxy hint so the cache can fetch upstream on miss.
* This is the only place feed media, videos and profile pictures are routed
* to the cache: callers keep the real URL, so the Tor decision, Coil/ExoPlayer
* cache keys and decryption-key lookups all see the origin.
*
* Coil's disk cache keys responses by the original `ImageRequest.data`, so
* disk caching continues to work transparently even though the network
* request now goes to localhost.
* Bridging is **opt-in**: only requests the media loaders mark with the
* [MEDIA_HEADER] header are considered (see [LocalBlossomMediaCallFactory]).
* Anything else whose last path segment merely looks like a hash — a BUD-02
* `GET /list/<pubkey>`, uploads, deletes, `HEAD` presence checks — must reach
* its server. The marker header is always stripped, so no server ever sees it;
* that is why Tor-proxied clients carry a non-bridging instance ([bridges] =
* false) instead of none: Tor-routed media skips the cache but is still cleaned.
*
* A marked `GET` is rewritten when [shouldBridge] returns `true`, it carries no
* `Authorization` (an auth-gated blob goes to the host the token was signed
* for), its last path segment is `<sha256>[.ext]`, and the host isn't already
* the cache. The original scheme+host(+path prefix) is passed as the `xs=`
* hint so the cache can fetch upstream on miss.
*
* Encrypted blobs are looked up in [keyCache] by URL after this interceptor
* runs, so their decryption key is registered under the rewritten URL too.
*
* When the cache refuses the connection (the app was closed since the last
* probe), [onUnreachable] is told so the bridge can switch off, and the
* request falls back to its original URL instead of failing.
*/
class LocalBlossomCacheRedirectInterceptor(
private val shouldBridge: () -> Boolean,
private val keyCache: EncryptionKeyCache? = null,
private val onUnreachable: () -> Unit = {},
val bridges: Boolean = true,
// Last, so the `LocalBlossomCacheRedirectInterceptor { enabled }` trailing-lambda form binds here.
// Told whether the request is a profile picture, for the "profile pictures only" setting.
private val shouldBridge: (profilePicture: Boolean) -> Boolean,
) : Interceptor {
override fun intercept(chain: Interceptor.Chain): Response {
val request = chain.request()
val marked = chain.request()
val kind = marked.header(MEDIA_HEADER)
val request = if (kind != null) marked.newBuilder().removeHeader(MEDIA_HEADER).build() else marked
if (!shouldBridge()) return chain.proceed(request)
if (request.method != "GET") return chain.proceed(request)
if (isLocalCache(request.url)) {
// Already addressed to the cache (a resolved `blossom:` URI): there is no
// original URL to fall back to here, but a dead cache must still be reported.
try {
return chain.proceed(request)
} catch (e: ConnectException) {
onUnreachable()
throw e
}
}
if (!bridges || kind == null || request.header("Authorization") != null) return chain.proceed(request)
if (!shouldBridge(kind == PROFILE_PICTURE)) return chain.proceed(request)
val rewritten = rewriteIfApplicable(request.url) ?: return chain.proceed(request)
return chain.proceed(
request
.newBuilder()
.url(rewritten)
.build(),
)
keyCache?.get(request.url.toString())?.let { keyCache.add(rewritten.toString(), it) }
val bridged =
try {
chain.proceed(
request
.newBuilder()
.url(rewritten)
.build(),
)
} catch (e: ConnectException) {
onUnreachable()
return chain.proceed(request)
}
if (bridged.isSuccessful) return bridged
// The cache answered, but not with the blob: it does not hold it and could not fetch it
// from `xs` either (not every cache implements that, and the one that does can be offline,
// still warming, or rate-limited). A miss is the ordinary state of a cache and must never
// be worse than having no cache at all, so the origin is asked directly — without this the
// 404 reached the caller and the image, video or encrypted file simply failed to load.
//
// Not reported through [onUnreachable]: the cache is alive and answering, so switching the
// bridge off would be the wrong conclusion.
bridged.close()
return chain.proceed(request)
}
private fun isLocalCache(url: HttpUrl): Boolean = url.port == LOCAL_CACHE_PORT && (url.host == LOCAL_CACHE_HOST || url.host.equals("localhost", ignoreCase = true))
private fun rewriteIfApplicable(url: HttpUrl): HttpUrl? {
val host = url.host
if (host == LOCAL_CACHE_HOST || host.equals("localhost", ignoreCase = true)) return null
@@ -121,5 +185,26 @@ class LocalBlossomCacheRedirectInterceptor(
const val LOCAL_CACHE_PORT = 24242
const val LOCAL_CACHE_BASE = "http://$LOCAL_CACHE_HOST:$LOCAL_CACHE_PORT"
private val BLOSSOM_LAST_SEGMENT_REGEX = Regex("^([0-9a-fA-F]{64})(?:\\.[^./]+)?$")
/** Marks a request as a media download the local cache may serve. Stripped before sending. */
const val MEDIA_HEADER = "X-Amethyst-Local-Blossom"
const val MEDIA = "media"
const val PROFILE_PICTURE = "profile-picture"
}
}
/**
* Marks every call it creates as a media download ([LocalBlossomCacheRedirectInterceptor.MEDIA_HEADER]),
* keeping a marker the request already carries (e.g. a profile picture's).
*/
class LocalBlossomMediaCallFactory(
private val delegate: Call.Factory,
private val kind: String = LocalBlossomCacheRedirectInterceptor.MEDIA,
) : Call.Factory {
override fun newCall(request: Request): Call =
if (request.header(LocalBlossomCacheRedirectInterceptor.MEDIA_HEADER) != null) {
delegate.newCall(request)
} else {
delegate.newCall(request.newBuilder().header(LocalBlossomCacheRedirectInterceptor.MEDIA_HEADER, kind).build())
}
}
@@ -55,12 +55,13 @@ class OkHttpClientFactory(
val userAgent: String,
private val dns: SurgeDns,
/**
* Returns `true` when sha256-keyed HTTP requests should be transparently
* rewritten to the local Blossom cache (master toggle on, profile-pictures-only
* restriction off, probe up). When `null`, the interceptor is disabled —
* useful for tests or pre-configuration call sites.
* Returns `true` when a sha256-keyed HTTP request should be transparently
* rewritten to the local Blossom cache, given whether it is a profile picture
* (master toggle on, probe up, and either not restricted to profile pictures
* or this is one). When `null`, the interceptor is disabled — useful for
* tests or pre-configuration call sites.
*/
val shouldBridgeBlossomCache: (() -> Boolean)? = null,
val shouldBridgeBlossomCache: ((profilePicture: Boolean) -> Boolean)? = null,
private val onionCache: OnionLocationCache,
/**
* Resource-usage ledger counter, installed OUTERMOST on the shared base
@@ -76,11 +77,24 @@ class OkHttpClientFactory(
* anonymous. See [BlossomReadAuthInterceptor].
*/
private val blossomReadAuth: Interceptor? = null,
/**
* Called when the local Blossom cache refuses a connection, so the caller can
* mark it unavailable instead of waiting for the next periodic probe.
*/
onLocalBlossomCacheUnreachable: () -> Unit = {},
) {
// val logging = LoggingInterceptor()
val keyDecryptor = EncryptedBlobInterceptor(keyCache)
// Tor-routed media skips the local cache (the cache would fetch the origin outside Tor, and
// Tor refuses 127.0.0.1 anyway), but its media marker header must still be stripped: this
// non-bridging instance takes the redirect's place on Tor clients, and stands in for it when
// the bridge isn't configured at all.
private val blossomCacheMarkerStripper = LocalBlossomCacheRedirectInterceptor(bridges = false) { false }
private val blossomCacheRedirect =
shouldBridgeBlossomCache?.let { LocalBlossomCacheRedirectInterceptor(it) }
shouldBridgeBlossomCache?.let { LocalBlossomCacheRedirectInterceptor(keyCache, onLocalBlossomCacheUnreachable, bridges = true, shouldBridge = it) }
?: blossomCacheMarkerStripper
// Most images/videos in a feed come from a small set of hosts (e.g. a single
// Blossom/imgproxy server). OkHttp's default dispatcher caps inflight requests
@@ -130,15 +144,13 @@ class OkHttpClientFactory(
.followSslRedirects(true)
.apply { usageInterceptor?.let { addInterceptor(it) } }
.addInterceptor(DefaultContentTypeInterceptor(userAgent))
.apply {
blossomCacheRedirect?.let { addInterceptor(it) }
}
// Sits outside the network interceptors so its retry re-runs the
// full stack (content-type, blossom cache, key decryptor) for the
// authenticated response. Only signs on an actual 401.
// Sits before the local-cache redirect and outside the network interceptors, so it
// sees the origin host (a token is scoped to it) and its retry re-runs the rest of
// the stack (blossom cache, key decryptor) for the authenticated response. Only
// signs on an actual 401. A request carrying a token is never bridged.
.apply {
blossomReadAuth?.let { addInterceptor(it) }
}
}.addInterceptor(blossomCacheRedirect)
// .addNetworkInterceptor(logging)
.addNetworkInterceptor(keyDecryptor)
// Passively populates [onionCache] from any HTTP/WebSocket response
@@ -172,7 +184,15 @@ class OkHttpClientFactory(
// `.onion`s — clearnet clients must never try to resolve `.onion`
// (DNS would fail, and we don't want fingerprintable lookups).
.apply { if (proxy != null) addInterceptor(OnionUrlRewriteInterceptor(onionCache)) }
.connectTimeout(Duration.ofSeconds(seconds.toLong()))
// The local Blossom cache lives on 127.0.0.1, which Tor refuses to reach, and it
// would fetch the origin outside Tor. Tor-routed media keeps going to its origin
// through Tor; the local cache is used by the direct client only.
.apply {
if (proxy != null) {
val index = interceptors().indexOf(blossomCacheRedirect)
if (index >= 0) interceptors()[index] = blossomCacheMarkerStripper
}
}.connectTimeout(Duration.ofSeconds(seconds.toLong()))
.readTimeout(Duration.ofSeconds(seconds.toLong() * 3))
.writeTimeout(Duration.ofSeconds(seconds.toLong() * 3))
.build()
@@ -0,0 +1,295 @@
/*
* Copyright (c) 2025 Vitor Pamplona
*
* Permission is hereby granted, free of charge, to any person obtaining a copy of
* this software and associated documentation files (the "Software"), to deal in
* the Software without restriction, including without limitation the rights to use,
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
* Software, and to permit persons to whom the Software is furnished to do so,
* subject to the following conditions:
*
* The above copyright notice and this permission notice shall be included in all
* copies or substantial portions of the Software.
*
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
*/
package com.vitorpamplona.amethyst.commons.service.http
import com.vitorpamplona.quartz.utils.ciphers.NostrCipher
import okhttp3.Interceptor
import okhttp3.OkHttpClient
import okhttp3.Protocol
import okhttp3.Request
import okhttp3.RequestBody.Companion.toRequestBody
import okhttp3.Response
import okhttp3.ResponseBody.Companion.toResponseBody
import java.lang.reflect.Proxy
import java.net.ConnectException
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFailsWith
import kotlin.test.assertNotNull
import kotlin.test.assertNull
import kotlin.test.assertSame
import kotlin.test.assertTrue
class LocalBlossomCacheRedirectSafetyTest {
private val sha = "b1674191a88ec5cdd733e4240a81803105dc412d6c6708d53ab94fc248f4f553"
private val origin = "https://blossom.example.com/$sha.jpg"
private val bridged = "http://127.0.0.1:24242/$sha.jpg?xs=https%3A%2F%2Fblossom.example.com"
private fun media(
url: String = origin,
kind: String = LocalBlossomCacheRedirectInterceptor.MEDIA,
) = Request.Builder().url(url).header(LocalBlossomCacheRedirectInterceptor.MEDIA_HEADER, kind)
private object IdentityCipher : NostrCipher {
override fun name() = "identity"
override fun encrypt(bytesToEncrypt: ByteArray) = bytesToEncrypt
override fun decrypt(bytesToDecrypt: ByteArray) = bytesToDecrypt
override fun decryptOrNull(bytesToDecrypt: ByteArray) = bytesToDecrypt
}
/** Records every request passed to proceed; [refuse] makes matching ones fail like a closed port. */
private fun chain(
request: Request,
sent: MutableList<Request>,
statusFor: (Request) -> Int = { 200 },
refuse: (Request) -> Boolean = { false },
): Interceptor.Chain =
Proxy.newProxyInstance(
Interceptor.Chain::class.java.classLoader,
arrayOf(Interceptor.Chain::class.java),
) { _, method, args ->
when (method.name) {
"request" -> request
"proceed" -> {
val proceeded = args[0] as Request
sent.add(proceeded)
if (refuse(proceeded)) throw ConnectException("Connection refused")
val status = statusFor(proceeded)
Response
.Builder()
.request(proceeded)
.protocol(Protocol.HTTP_1_1)
.code(status)
.message(if (status == 200) "OK" else "Not Found")
.body("".toResponseBody(null))
.build()
}
else -> throw UnsupportedOperationException(method.name)
}
} as Interceptor.Chain
@Test
fun getIsBridged() {
val sent = mutableListOf<Request>()
LocalBlossomCacheRedirectInterceptor { true }.intercept(chain(media().build(), sent)).close()
assertEquals(bridged, sent.single().url.toString())
}
@Test
fun profilePicturesOnlyBridgesProfilePicturesAlone() {
// The "profile pictures only" setting: bridge a request only when it is a profile picture.
val interceptor = LocalBlossomCacheRedirectInterceptor { profilePicture -> profilePicture }
val feedImage = mutableListOf<Request>()
interceptor.intercept(chain(media().build(), feedImage)).close()
assertEquals(origin, feedImage.single().url.toString())
val avatar = mutableListOf<Request>()
interceptor.intercept(chain(media(kind = LocalBlossomCacheRedirectInterceptor.PROFILE_PICTURE).build(), avatar)).close()
assertEquals(bridged, avatar.single().url.toString())
}
@Test
fun unmarkedRequestsAreNeverBridged() {
// BUD-02 `GET /list/<pubkey>` ends in 64 hex chars too, but it is not a blob download.
val list = "https://blossom.example.com/list/$sha"
val sent = mutableListOf<Request>()
LocalBlossomCacheRedirectInterceptor { true }
.intercept(
chain(
Request
.Builder()
.url(list)
.header("Authorization", "Nostr abc")
.build(),
sent,
),
).close()
assertEquals(list, sent.single().url.toString())
}
@Test
fun authGatedDownloadsGoToTheirServer() {
val sent = mutableListOf<Request>()
LocalBlossomCacheRedirectInterceptor { true }
.intercept(chain(media().header("Authorization", "Nostr abc").build(), sent))
.close()
assertEquals(origin, sent.single().url.toString())
}
@Test
fun theMarkerNeverLeavesTheApp() {
listOf(
LocalBlossomCacheRedirectInterceptor { true },
LocalBlossomCacheRedirectInterceptor { false },
LocalBlossomCacheRedirectInterceptor(bridges = false) { true },
).forEach { interceptor ->
val sent = mutableListOf<Request>()
interceptor.intercept(chain(media().build(), sent)).close()
assertNull(sent.single().header(LocalBlossomCacheRedirectInterceptor.MEDIA_HEADER))
}
}
@Test
fun aNonBridgingInstanceLeavesMediaAlone() {
val sent = mutableListOf<Request>()
LocalBlossomCacheRedirectInterceptor(bridges = false) { true }.intercept(chain(media().build(), sent)).close()
assertEquals(origin, sent.single().url.toString())
}
@Test
fun mediaCallFactoryMarksCallsButKeepsAnExistingMarker() {
val seen = mutableListOf<Request>()
val factory =
LocalBlossomMediaCallFactory({ request ->
seen.add(request)
OkHttpClient().newCall(request)
})
factory.newCall(Request.Builder().url(origin).build())
factory.newCall(media(kind = LocalBlossomCacheRedirectInterceptor.PROFILE_PICTURE).build())
assertEquals(
listOf(LocalBlossomCacheRedirectInterceptor.MEDIA, LocalBlossomCacheRedirectInterceptor.PROFILE_PICTURE),
seen.map { it.header(LocalBlossomCacheRedirectInterceptor.MEDIA_HEADER) },
)
}
@Test
fun serverSpecificMethodsAreNotBridged() {
val interceptor = LocalBlossomCacheRedirectInterceptor { true }
val requests =
listOf(
Request
.Builder()
.url(origin)
.head()
.build(),
Request
.Builder()
.url(origin)
.delete()
.header("Authorization", "Nostr abc")
.build(),
Request
.Builder()
.url(origin)
.put("x".toRequestBody())
.build(),
)
requests.forEach { request ->
val sent = mutableListOf<Request>()
interceptor.intercept(chain(request, sent)).close()
assertEquals(origin, sent.single().url.toString(), "${request.method} must reach its server")
}
}
@Test
fun decryptionKeyFollowsTheRewrite() {
val keys = EncryptionKeyCache()
val info = DecryptInformation(IdentityCipher, "image/jpeg")
keys.add(origin, info)
val sent = mutableListOf<Request>()
LocalBlossomCacheRedirectInterceptor(keyCache = keys) { true }.intercept(chain(media().build(), sent)).close()
// EncryptedBlobInterceptor runs after the rewrite and looks the key up by the URL it sees.
assertSame(info, assertNotNull(keys.get(sent.single().url.toString())))
}
@Test
fun deadCacheFallsBackToOriginAndIsReported() {
var reports = 0
val sent = mutableListOf<Request>()
val response =
LocalBlossomCacheRedirectInterceptor(onUnreachable = { reports++ }) { true }
.intercept(chain(media().build(), sent) { it.url.host == "127.0.0.1" })
assertEquals(200, response.code)
assertEquals(listOf(bridged, origin), sent.map { it.url.toString() })
assertEquals(1, reports)
response.close()
}
@Test
fun deadCacheIsReportedForResolvedBlossomUrls() {
var reports = 0
val sent = mutableListOf<Request>()
val request = media("http://127.0.0.1:24242/$sha.mp4?xs=https://blossom.example.com").build()
assertFailsWith<ConnectException> {
LocalBlossomCacheRedirectInterceptor(onUnreachable = { reports++ }) { true }
.intercept(chain(request, sent) { true })
}
assertEquals(1, reports)
assertTrue(sent.single().url.host == "127.0.0.1")
}
@Test
fun aCacheMissFallsBackToTheOrigin() {
var reports = 0
val sent = mutableListOf<Request>()
val response =
LocalBlossomCacheRedirectInterceptor(onUnreachable = { reports++ }) { true }
.intercept(
chain(media().build(), sent, statusFor = { if (it.url.host == "127.0.0.1") 404 else 200 }),
)
// The blob still arrives: the cache is an optimisation, never a gate.
assertEquals(200, response.code)
assertEquals(origin, response.request.url.toString())
assertEquals(listOf(bridged, origin), sent.map { it.url.toString() })
// The cache answered, so it is alive — the bridge must not be switched off.
assertEquals(0, reports)
response.close()
}
@Test
fun aCacheErrorFallsBackToTheOriginToo() {
val sent = mutableListOf<Request>()
val response =
LocalBlossomCacheRedirectInterceptor { true }
.intercept(
chain(media().build(), sent, statusFor = { if (it.url.host == "127.0.0.1") 500 else 200 }),
)
assertEquals(200, response.code)
assertEquals(listOf(bridged, origin), sent.map { it.url.toString() })
response.close()
}
@Test
fun aServedBlobIsNotRefetchedFromTheOrigin() {
val sent = mutableListOf<Request>()
val response = LocalBlossomCacheRedirectInterceptor { true }.intercept(chain(media().build(), sent))
assertEquals(200, response.code)
// A hit must cost exactly one request, to the cache.
assertEquals(listOf(bridged), sent.map { it.url.toString() })
response.close()
}
}
@@ -0,0 +1,115 @@
/*
* Copyright (c) 2025 Vitor Pamplona
*
* Permission is hereby granted, free of charge, to any person obtaining a copy of
* this software and associated documentation files (the "Software"), to deal in
* the Software without restriction, including without limitation the rights to use,
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
* Software, and to permit persons to whom the Software is furnished to do so,
* subject to the following conditions:
*
* The above copyright notice and this permission notice shall be included in all
* copies or substantial portions of the Software.
*
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
*/
package com.vitorpamplona.amethyst.commons.service.http
import com.vitorpamplona.quartz.nip01Core.relay.sockets.okhttp.SurgeDns
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.cancel
import kotlinx.coroutines.flow.MutableStateFlow
import okhttp3.Request
import java.net.ServerSocket
import kotlin.concurrent.thread
import kotlin.test.AfterTest
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertFalse
import kotlin.test.assertTrue
/**
* With "videos via Tor" on, the player pool is built with a Tor-proxied [DynamicCallFactory], but a
* `blossom:` video resolves to the local Blossom cache on 127.0.0.1 only when the data source opens.
* Tor refuses loopback addresses, so the local cache must be reached directly or no video plays.
*/
class LocalBlossomCacheTorRoutingTest {
private val scope = CoroutineScope(SupervisorJob() + Dispatchers.Default)
@AfterTest
fun tearDown() = scope.cancel()
private fun manager(proxyPort: Int) =
DualHttpClientManager(
userAgent = "test",
proxyPortProvider = MutableStateFlow(proxyPort),
isMobileDataProvider = MutableStateFlow(false),
keyCache = EncryptionKeyCache(),
scope = scope,
dns = SurgeDns(),
shouldBridgeBlossomCache = { true },
onionCache = OnionLocationCache(),
)
/** A port nothing listens on: stands in for a SOCKS proxy that cannot reach the target. */
private fun deadPort(): Int = ServerSocket(0).use { it.localPort }
/** Minimal one-shot HTTP server standing in for the local Blossom cache. */
private fun localCache(body: String): ServerSocket {
val server = ServerSocket(0)
thread(isDaemon = true) {
server.accept().use { socket ->
// Drain the request headers before answering.
val reader = socket.getInputStream().bufferedReader()
do {
val line = reader.readLine()
} while (!line.isNullOrEmpty())
socket.getOutputStream().write(
"HTTP/1.1 200 OK\r\nContent-Length: ${body.length}\r\nConnection: close\r\n\r\n$body".toByteArray(),
)
}
}
return server
}
@Test
fun proxiedFactoryReachesLoopbackCacheDirectly() {
val cache = localCache("video-bytes")
val factory = manager(deadPort()).getDynamicCallFactory(useProxy = true)
val sha = "b1674191a88ec5cdd733e4240a81803105dc412d6c6708d53ab94fc248f4f553"
val request = Request.Builder().url("http://127.0.0.1:${cache.localPort}/$sha.mp4?xs=https://blossom.example.com").build()
factory.newCall(request).execute().use { response ->
assertEquals(200, response.code)
assertEquals("video-bytes", response.body.string())
}
cache.close()
}
@Test
fun proxyDecisionUsesTheFinalRequestUrl() {
assertFalse(DynamicCallFactory.shouldUseProxy(true, "http://127.0.0.1:24242/abc.mp4?xs=https://blossom.primal.net"))
assertFalse(DynamicCallFactory.shouldUseProxy(true, "http://localhost:24242/abc.mp4"))
assertTrue(DynamicCallFactory.shouldUseProxy(true, "https://blossom.primal.net/abc.mp4"))
assertFalse(DynamicCallFactory.shouldUseProxy(false, "https://blossom.primal.net/abc.mp4"))
}
@Test
fun onlyTheDirectClientRewritesToTheLocalCache() {
val manager = manager(deadPort())
val direct = manager.getHttpClient(useProxy = false)
val proxied = manager.getHttpClient(useProxy = true)
// Both carry one (it strips the media marker), but only the direct client's bridges.
assertEquals(listOf(true), direct.interceptors.filterIsInstance<LocalBlossomCacheRedirectInterceptor>().map { it.bridges })
assertEquals(listOf(false), proxied.interceptors.filterIsInstance<LocalBlossomCacheRedirectInterceptor>().map { it.bridges })
}
}