diff --git a/Cargo.lock b/Cargo.lock index ee5d66fa6..c4a5996c5 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1785,8 +1785,10 @@ name = "fluxer-gifs" version = "0.1.0" dependencies = [ "anyhow", + "axum", "fluxer-svc", "fluxer_common", + "futures", "hmac 0.13.0", "moka", "reqwest", diff --git a/fluxer_gifs/Cargo.toml b/fluxer_gifs/Cargo.toml index 55f02fb30..d84d94225 100644 --- a/fluxer_gifs/Cargo.toml +++ b/fluxer_gifs/Cargo.toml @@ -11,6 +11,7 @@ publish = false fluxer_common = { path = "../fluxer_common" } anyhow = "1.0.104" fluxer-svc = { path = "../fluxer_svc", default-features = false } +futures = "0.3" hmac = "0.13.0" moka = { version = "0.12.16", features = ["future", "sync"] } reqwest = { version = "0.13.5", default-features = false, features = ["json", "rustls"] } @@ -21,3 +22,7 @@ tokio = { version = "1.53.1", features = ["macros", "rt-multi-thread", "signal", tracing = "0.1.44" url = "2.5" urlencoding = "2.1" + +[dev-dependencies] +axum = { version = "0.8.9", default-features = false, features = ["http1", "json", "query", "tokio"] } +tokio = { version = "1.53.1", features = ["net"] } diff --git a/fluxer_gifs/src/klipy.rs b/fluxer_gifs/src/klipy.rs index ab3c2b9f9..de29d3115 100644 --- a/fluxer_gifs/src/klipy.rs +++ b/fluxer_gifs/src/klipy.rs @@ -3,6 +3,7 @@ use crate::media_proxy::MediaProxyUrlBuilder; use crate::types::{GifCategoryTag, GifItem, GifMediaFormat}; use anyhow::Context; +use futures::stream::{self, StreamExt}; use reqwest::Url; use serde::Deserialize; use serde_json::Value; @@ -19,6 +20,7 @@ const MAX_RETRIES: usize = 3; const BACKOFF_BASE_DELAY: Duration = Duration::from_secs(1); const KLIPY_RESPONSE_LIMIT_BYTES: usize = 512 * 1024; const MAX_FEATURED_CATEGORIES: usize = 50; +const FEATURED_CATEGORY_PREVIEW_CONCURRENCY: usize = 10; const FEATURED_CATEGORIES_FETCH_TIMEOUT: Duration = Duration::from_secs(30); const FLUXER_USER_AGENT: &str = "Fluxerbot/1.0 (+https://fluxer.app)"; const KLIPY_PROVIDER_NAME: &str = "klipy"; @@ -54,6 +56,7 @@ static MEDIA_FILTER: LazyLock = LazyLock::new(|| MEDIA_FORMAT_PREFERENCE pub struct KlipyClient { http_client: reqwest::Client, media_proxy: MediaProxyUrlBuilder, + base_url: String, } #[derive(Debug, Deserialize)] @@ -158,9 +161,16 @@ impl KlipyClient { Ok(Self { http_client, media_proxy, + base_url: KLIPY_BASE_URL.to_owned(), }) } + #[cfg(test)] + fn with_base_url(mut self, base_url: &str) -> Self { + self.base_url = base_url.trim_end_matches('/').to_owned(); + self + } + pub async fn search( &self, api_key: &str, @@ -335,31 +345,37 @@ impl KlipyClient { .take(MAX_FEATURED_CATEGORIES) .collect::>(); - let mut categories = Vec::with_capacity(search_terms.len()); - for search_term in search_terms { - let gif = match self - .search( - api_key, - &search_term, - &normalized_locale, - KLIPY_FEATURED_CATEGORY_REFRESH_COUNTRY, - 1, - ) - .await - { - Ok(mut gifs) => gifs.drain(..).next(), - Err(err) => { - tracing::debug!( - error = %err, - search_term = %search_term, - locale = %normalized_locale, - "failed to fetch KLIPY category preview GIF" - ); - None + let categories = stream::iter(search_terms) + .map(|search_term| { + let normalized_locale = &normalized_locale; + async move { + let gif = match self + .search( + api_key, + &search_term, + normalized_locale, + KLIPY_FEATURED_CATEGORY_REFRESH_COUNTRY, + 1, + ) + .await + { + Ok(mut gifs) => gifs.drain(..).next(), + Err(err) => { + tracing::debug!( + error = %err, + search_term = %search_term, + locale = %normalized_locale, + "failed to fetch KLIPY category preview GIF" + ); + None + } + }; + category_response(search_term, gif) } - }; - categories.push(category_response(search_term, gif)); - } + }) + .buffered(FEATURED_CATEGORY_PREVIEW_CONCURRENCY) + .collect::>() + .await; Ok(categories) } @@ -498,7 +514,7 @@ impl KlipyClient { } fn create_url(&self, endpoint: &str, params: &[(&str, &str)]) -> anyhow::Result { - let mut url = Url::parse(&format!("{KLIPY_BASE_URL}/{endpoint}"))?; + let mut url = Url::parse(&format!("{}/{endpoint}", self.base_url))?; { let mut query = url.query_pairs_mut(); query.append_pair("client_key", CLIENT_KEY); @@ -841,6 +857,72 @@ fn category_response(name: String, gif: Option) -> GifCategoryTag { mod tests { use super::*; + #[tokio::test(flavor = "multi_thread")] + async fn cold_featured_categories_fit_the_router_shard_budget() { + use axum::extract::Query; + use axum::routing::get; + use std::collections::HashMap; + + const UPSTREAM_LATENCY: Duration = Duration::from_millis(200); + + let categories = || async { + sleep(UPSTREAM_LATENCY).await; + let tags = (0..MAX_FEATURED_CATEGORIES) + .map(|index| serde_json::json!({"searchterm": format!("term-{index}")})) + .collect::>(); + axum::Json(serde_json::json!({ "tags": tags })) + }; + let search = |Query(params): Query>| async move { + sleep(UPSTREAM_LATENCY).await; + let term = params.get("q").cloned().unwrap_or_default(); + axum::Json(serde_json::json!({ + "results": [{ + "id": term, + "title": term, + "itemurl": format!("https://klipy.com/gifs/{term}"), + "media_formats": { + "gif": {"url": format!("https://static.klipy.com/{term}.gif"), "dims": [100, 100]} + } + }] + })) + }; + let app = axum::Router::new() + .route("/categories", get(categories)) + .route("/search", get(search)); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0") + .await + .expect("bind"); + let address = listener.local_addr().expect("local addr"); + tokio::spawn(async move { axum::serve(listener, app).await }); + + let client = KlipyClient::new(MediaProxyUrlBuilder::for_test( + "https://media.example.test", + "secret", + )) + .expect("client") + .with_base_url(&format!("http://{address}")); + + let started = std::time::Instant::now(); + let categories = client + .featured_categories("key", "sv-SE") + .await + .expect("featured categories"); + let elapsed = started.elapsed(); + + assert_eq!(MAX_FEATURED_CATEGORIES, categories.len()); + assert!(categories.iter().all(|category| !category.src.is_empty())); + assert_eq!("term-0", categories[0].name); + assert_eq!( + format!("term-{}", MAX_FEATURED_CATEGORIES - 1), + categories[MAX_FEATURED_CATEGORIES - 1].name + ); + assert!( + elapsed < fluxer_svc::router::SHARD_REQUEST_TIMEOUT / 2, + "cold featured categories took {elapsed:?}, the router gives the shard {:?}", + fluxer_svc::router::SHARD_REQUEST_TIMEOUT + ); + } + #[test] fn locale_uses_klipy_supported_form() { assert_eq!(normalize_locale("en-US"), "en_US"); diff --git a/fluxer_svc/src/router.rs b/fluxer_svc/src/router.rs index 3a551e963..712612444 100644 --- a/fluxer_svc/src/router.rs +++ b/fluxer_svc/src/router.rs @@ -17,7 +17,7 @@ use tokio::sync::{Semaphore, TryAcquireError}; use tokio::task::JoinSet; use tracing::{debug, info, warn}; -pub(crate) const SHARD_REQUEST_TIMEOUT: Duration = Duration::from_secs(5); +pub const SHARD_REQUEST_TIMEOUT: Duration = Duration::from_secs(5); const MAX_BROADCAST_CONCURRENCY: usize = 32; const MAX_ROUTER_REQUEST_BYTES: usize = 2 * 1024 * 1024; const LEGACY_SHARD_DECODE_ERROR: &[u8] = br#"{"error":"shard_request_decode_error"}"#;