fix: migrate gifs to klipy service (#1220)

This commit is contained in:
Hampus
2026-06-30 01:19:57 +02:00
committed by GitHub
parent 99eedf6d1e
commit 4a741d8396
79 changed files with 3032 additions and 1653 deletions
+22
View File
@@ -0,0 +1,22 @@
# SPDX-License-Identifier: AGPL-3.0-or-later
[package]
name = "fluxer-gifs"
version = "0.1.0"
edition.workspace = true
license.workspace = true
[dependencies]
anyhow = "1.0.102"
base64 = "0.22.1"
fluxer-svc = { path = "../fluxer_svc", default-features = false }
hmac = "0.13.0"
moka = { version = "0.12.15", features = ["future", "sync"] }
reqwest = { version = "0.13.4", default-features = false, features = ["json", "rustls"] }
serde = { version = "1.0.228", features = ["derive"] }
serde_json = "1.0.150"
sha2 = "0.11.0"
tokio = { version = "1.52.3", features = ["macros", "rt-multi-thread", "signal", "sync", "time"] }
tracing = "0.1.44"
url = "2.5"
urlencoding = "2.1"
+28
View File
@@ -0,0 +1,28 @@
# SPDX-License-Identifier: AGPL-3.0-or-later
FROM rust:1-bookworm AS builder
WORKDIR /usr/src/app
COPY . .
RUN cargo build --release -p fluxer-gifs
FROM debian:bookworm-slim
ARG BUILD_VERSION=""
WORKDIR /usr/local/bin
RUN apt-get update && apt-get install -y --no-install-recommends ca-certificates && \
rm -rf /var/lib/apt/lists/*
COPY --from=builder /usr/src/app/target/release/fluxer-gifs /usr/local/bin/fluxer-gifs
ENV BUILD_VERSION="${BUILD_VERSION}"
USER 65532:65532
EXPOSE 8090
CMD ["/usr/local/bin/fluxer-gifs"]
+721
View File
@@ -0,0 +1,721 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
use crate::media_proxy::MediaProxyUrlBuilder;
use crate::types::{GifCategoryTag, GifItem, GifMediaFormat};
use anyhow::Context;
use reqwest::Url;
use serde::Deserialize;
use serde_json::Value;
use std::collections::{BTreeMap, HashSet};
use std::time::Duration;
use tokio::time::sleep;
const KLIPY_BASE_URL: &str = "https://api.klipy.com/v2";
const KLIPY_DIRECT_BASE_URL: &str = "https://api.klipy.com/api/v1";
const DEFAULT_CONTENT_FILTER: &str = "low";
const CLIENT_KEY: &str = "fluxer";
const MAX_RETRIES: usize = 3;
const BACKOFF_BASE_DELAY: Duration = Duration::from_secs(1);
const KLIPY_RESPONSE_LIMIT_BYTES: usize = 512 * 1024;
const FLUXER_USER_AGENT: &str = "Fluxerbot/1.0 (+https://fluxer.app)";
const KLIPY_PROVIDER_NAME: &str = "klipy";
const KLIPY_FEATURED_CATEGORY_REFRESH_COUNTRY: &str = "US";
const SIZE_PREFERENCE: [&str; 4] = ["hd", "md", "sm", "xs"];
const FORMAT_PREFERENCE: [&str; 4] = ["webm", "mp4", "webp", "gif"];
#[derive(Clone)]
pub struct KlipyClient {
http_client: reqwest::Client,
media_proxy: MediaProxyUrlBuilder,
}
#[derive(Debug, Deserialize)]
struct ResultsResponse {
results: Vec<Value>,
}
#[derive(Debug, Deserialize)]
struct TagsResponse {
tags: Vec<Value>,
}
#[derive(Debug, Deserialize)]
struct KlipyGif {
id: Value,
#[serde(default)]
slug: Option<String>,
#[serde(default)]
title: String,
#[serde(default)]
itemurl: Option<String>,
#[serde(default)]
file: Option<BTreeMap<String, BTreeMap<String, KlipyFileEntry>>>,
#[serde(default)]
media_formats: Option<KlipyMediaFormats>,
}
#[derive(Debug, Deserialize)]
struct KlipyMediaFormats {
#[serde(default)]
webm: Option<KlipyFallbackMediaFormat>,
}
#[derive(Debug, Deserialize)]
struct KlipyFallbackMediaFormat {
url: String,
dims: [i32; 2],
}
#[derive(Debug, Deserialize)]
struct KlipyFileEntry {
#[serde(default)]
url: Option<String>,
#[serde(default)]
width: Option<i32>,
#[serde(default)]
height: Option<i32>,
}
#[derive(Debug, Deserialize)]
struct KlipyCategoryTag {
searchterm: String,
}
#[derive(Debug, Deserialize)]
struct DirectGifResponse {
#[serde(default)]
data: Option<KlipyGif>,
}
enum KlipyJsonFetch<T> {
Found(T),
NotFound,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct KlipyPath {
path_type: KlipyPathType,
slug: String,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum KlipyPathType {
Gif,
Clip,
}
impl KlipyClient {
pub fn new(media_proxy: MediaProxyUrlBuilder) -> anyhow::Result<Self> {
let http_client = reqwest::Client::builder()
.user_agent(FLUXER_USER_AGENT)
.timeout(Duration::from_secs(30))
.build()
.context("failed to build KLIPY HTTP client")?;
Ok(Self {
http_client,
media_proxy,
})
}
pub async fn search(
&self,
api_key: &str,
q: &str,
locale: &str,
country: &str,
limit: u32,
) -> anyhow::Result<Vec<GifItem>> {
let locale = normalize_locale(locale);
let limit = limit.to_string();
self.fetch_gifs(
"search",
&[
("key", api_key),
("q", q),
("country", country),
("locale", &locale),
("limit", &limit),
],
)
.await
}
pub async fn featured_gifs(
&self,
api_key: &str,
locale: &str,
country: &str,
) -> anyhow::Result<Vec<GifItem>> {
let locale = normalize_locale(locale);
self.fetch_gifs(
"featured",
&[
("key", api_key),
("country", country),
("locale", &locale),
("limit", "1"),
],
)
.await
}
pub async fn trending_gifs(
&self,
api_key: &str,
locale: &str,
country: &str,
) -> anyhow::Result<Vec<GifItem>> {
let locale = normalize_locale(locale);
self.fetch_gifs(
"featured",
&[
("key", api_key),
("country", country),
("locale", &locale),
("limit", "50"),
],
)
.await
}
pub async fn suggestions(
&self,
api_key: &str,
q: &str,
locale: &str,
) -> anyhow::Result<Vec<String>> {
let locale = normalize_locale(locale);
let response: ResultsResponse = self
.fetch_json(
"autocomplete",
&[("key", api_key), ("q", q), ("locale", &locale)],
)
.await?;
Ok(response
.results
.into_iter()
.filter_map(|value| value.as_str().map(ToOwned::to_owned))
.collect())
}
pub async fn register_share(
&self,
api_key: &str,
id: &str,
q: &str,
locale: &str,
country: &str,
) -> anyhow::Result<()> {
let locale = normalize_locale(locale);
let url = self.create_url(
"registershare",
&[
("key", api_key),
("id", id),
("country", country),
("locale", &locale),
("q", q),
],
)?;
let response = self.http_client.get(url).send().await?;
if !response.status().is_success() {
anyhow::bail!(
"KLIPY registershare failed with status {}",
response.status()
);
}
Ok(())
}
pub async fn resolve_by_url(
&self,
api_key: &str,
url: &str,
_locale: &str,
_country: &str,
) -> anyhow::Result<Option<GifItem>> {
let Some(path) = parse_klipy_path(url) else {
return Ok(None);
};
self.fetch_direct_gif(api_key, &path).await
}
pub async fn featured_categories(
&self,
api_key: &str,
locale: &str,
) -> anyhow::Result<Vec<GifCategoryTag>> {
let normalized_locale = normalize_locale(locale);
let response: TagsResponse = self
.fetch_json(
"categories",
&[
("key", api_key),
("country", KLIPY_FEATURED_CATEGORY_REFRESH_COUNTRY),
("locale", &normalized_locale),
("type", "featured"),
],
)
.await?;
let mut seen = HashSet::new();
let search_terms = response
.tags
.into_iter()
.filter_map(|value| serde_json::from_value::<KlipyCategoryTag>(value).ok())
.map(|tag| tag.searchterm.trim().to_owned())
.filter(|term| !term.is_empty())
.filter(|term| seen.insert(term.clone()))
.collect::<Vec<_>>();
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
}
};
categories.push(category_response(search_term, gif));
}
Ok(categories)
}
async fn fetch_gifs(
&self,
endpoint: &str,
params: &[(&str, &str)],
) -> anyhow::Result<Vec<GifItem>> {
let response: ResultsResponse = self.fetch_json(endpoint, params).await?;
Ok(response
.results
.into_iter()
.filter_map(|value| serde_json::from_value::<KlipyGif>(value).ok())
.filter_map(|gif| self.transform_gif(gif))
.collect())
}
async fn fetch_direct_gif(
&self,
api_key: &str,
path: &KlipyPath,
) -> anyhow::Result<Option<GifItem>> {
let url = self.create_direct_url(api_key, path)?;
let mut last_error = None;
for attempt in 0..MAX_RETRIES {
match self.fetch_direct_gif_once(url.clone(), path).await {
Ok(value) => return Ok(value),
Err(error) if attempt + 1 < MAX_RETRIES => {
last_error = Some(error);
sleep(BACKOFF_BASE_DELAY * 2_u32.pow(attempt as u32)).await;
}
Err(error) => return Err(error),
}
}
Err(last_error.unwrap_or_else(|| anyhow::anyhow!("exceeded KLIPY retry limit")))
}
async fn fetch_direct_gif_once(
&self,
url: Url,
path: &KlipyPath,
) -> anyhow::Result<Option<GifItem>> {
match self.fetch_json_response::<DirectGifResponse>(url).await? {
KlipyJsonFetch::NotFound => Ok(None),
KlipyJsonFetch::Found(response) => Ok(response
.data
.and_then(|gif| self.transform_gif_with_path(gif, Some(path)))),
}
}
async fn fetch_json<T>(&self, endpoint: &str, params: &[(&str, &str)]) -> anyhow::Result<T>
where
T: serde::de::DeserializeOwned,
{
let url = self.create_url(endpoint, params)?;
let mut last_error = None;
for attempt in 0..MAX_RETRIES {
match self.fetch_json_once(url.clone()).await {
Ok(value) => return Ok(value),
Err(error) if attempt + 1 < MAX_RETRIES => {
last_error = Some(error);
sleep(BACKOFF_BASE_DELAY * 2_u32.pow(attempt as u32)).await;
}
Err(error) => return Err(error),
}
}
Err(last_error.unwrap_or_else(|| anyhow::anyhow!("exceeded KLIPY retry limit")))
}
async fn fetch_json_once<T>(&self, url: Url) -> anyhow::Result<T>
where
T: serde::de::DeserializeOwned,
{
match self.fetch_json_response(url).await? {
KlipyJsonFetch::Found(value) => Ok(value),
KlipyJsonFetch::NotFound => anyhow::bail!("KLIPY request returned not found"),
}
}
async fn fetch_json_response<T>(&self, url: Url) -> anyhow::Result<KlipyJsonFetch<T>>
where
T: serde::de::DeserializeOwned,
{
let response = self.http_client.get(url.clone()).send().await?;
if response.status() == reqwest::StatusCode::NOT_FOUND {
return Ok(KlipyJsonFetch::NotFound);
}
if !response.status().is_success() {
anyhow::bail!("KLIPY request failed with status {}", response.status());
}
if response
.content_length()
.is_some_and(|len| len > KLIPY_RESPONSE_LIMIT_BYTES as u64)
{
anyhow::bail!("KLIPY response declared more than {KLIPY_RESPONSE_LIMIT_BYTES} bytes");
}
let bytes = response.bytes().await?;
if bytes.len() > KLIPY_RESPONSE_LIMIT_BYTES {
anyhow::bail!("KLIPY response exceeded {KLIPY_RESPONSE_LIMIT_BYTES} bytes");
}
serde_json::from_slice(&bytes)
.with_context(|| format!("failed to parse KLIPY response from {url}"))
.map(KlipyJsonFetch::Found)
}
fn create_url(&self, endpoint: &str, params: &[(&str, &str)]) -> anyhow::Result<Url> {
let mut url = Url::parse(&format!("{KLIPY_BASE_URL}/{endpoint}"))?;
{
let mut query = url.query_pairs_mut();
query.append_pair("client_key", CLIENT_KEY);
query.append_pair("contentfilter", DEFAULT_CONTENT_FILTER);
for (key, value) in params {
query.append_pair(key, value);
}
}
Ok(url)
}
fn create_direct_url(&self, api_key: &str, path: &KlipyPath) -> anyhow::Result<Url> {
let mut url = Url::parse(&format!("{KLIPY_DIRECT_BASE_URL}/"))?;
{
let mut segments = url
.path_segments_mut()
.map_err(|_| anyhow::anyhow!("KLIPY direct base URL cannot be a base"))?;
segments
.push(api_key)
.push(klipy_resource(path.path_type))
.push(&path.slug);
}
Ok(url)
}
fn transform_gif(&self, input: KlipyGif) -> Option<GifItem> {
self.transform_gif_with_path(input, None)
}
fn transform_gif_with_path(
&self,
input: KlipyGif,
fallback_path: Option<&KlipyPath>,
) -> Option<GifItem> {
let parsed_path = input.itemurl.as_deref().and_then(parse_klipy_path);
let resolved_path = parsed_path.as_ref().or(fallback_path);
let explicit_slug = input
.slug
.as_deref()
.map(str::trim)
.filter(|slug| !slug.is_empty());
let fallback_id = klipy_id_as_string(&input.id)?;
let normalized_slug = explicit_slug
.or_else(|| resolved_path.map(|path| path.slug.as_str()))
.unwrap_or(fallback_id.as_str())
.to_owned();
let normalized_type = resolved_path
.map(|path| path.path_type)
.unwrap_or(KlipyPathType::Gif);
let normalized_url = if resolved_path.is_some() || explicit_slug.is_some() {
build_share_url_with_type(normalized_type, &normalized_slug)
} else {
input
.itemurl
.clone()
.unwrap_or_else(|| build_share_url_with_type(normalized_type, &normalized_slug))
};
let (media, preferred) = self.collect_media(&input);
let top = media.get("webm").cloned().or(preferred)?;
Some(GifItem {
id: normalized_slug.clone(),
slug: normalized_slug,
provider: KLIPY_PROVIDER_NAME.to_owned(),
title: input.title,
url: normalized_url,
src: top.src.clone(),
proxy_src: top.proxy_src.clone(),
width: top.width,
height: top.height,
media,
placeholder: None,
})
}
fn collect_media(
&self,
input: &KlipyGif,
) -> (BTreeMap<String, GifMediaFormat>, Option<GifMediaFormat>) {
let mut media = BTreeMap::new();
let mut preferred = None;
for size in SIZE_PREFERENCE {
let Some(bucket) = input.file.as_ref().and_then(|files| files.get(size)) else {
continue;
};
for format in FORMAT_PREFERENCE {
let Some(entry) = bucket.get(format) else {
continue;
};
let Some(media_format) = self.to_media_format(entry) else {
continue;
};
let public_key = public_format_key(size, format);
media.insert(public_key, media_format.clone());
if preferred.is_none() {
preferred = Some(media_format);
}
}
}
if media.is_empty()
&& let Some(webm) = input
.media_formats
.as_ref()
.and_then(|formats| formats.webm.as_ref())
&& webm.dims[0] > 0
&& webm.dims[1] > 0
&& let Some(proxy_src) = self.media_proxy.external_proxy_url(&webm.url)
{
let fallback = GifMediaFormat {
src: webm.url.clone(),
proxy_src,
width: webm.dims[0],
height: webm.dims[1],
};
media.insert("webm".to_owned(), fallback.clone());
preferred = Some(fallback);
}
(media, preferred)
}
fn to_media_format(&self, entry: &KlipyFileEntry) -> Option<GifMediaFormat> {
let src = entry.url.as_ref()?;
let width = entry.width.filter(|width| *width > 0)?;
let height = entry.height.filter(|height| *height > 0)?;
let proxy_src = self.media_proxy.external_proxy_url(src)?;
Some(GifMediaFormat {
src: src.clone(),
proxy_src,
width,
height,
})
}
}
pub fn normalize_locale(locale: &str) -> String {
locale.replace('-', "_")
}
pub fn build_share_url(slug: &str) -> String {
let trimmed = slug.trim();
if trimmed.is_empty() {
return "https://klipy.com/gifs".to_owned();
}
build_share_url_with_type(KlipyPathType::Gif, trimmed)
}
pub fn extract_slug_from_url(url: &str) -> Option<String> {
parse_klipy_path(url).map(|path| path.slug)
}
fn klipy_id_as_string(value: &Value) -> Option<String> {
match value {
Value::String(value) => {
let trimmed = value.trim();
(!trimmed.is_empty()).then(|| trimmed.to_owned())
}
Value::Number(value) => Some(value.to_string()),
_ => None,
}
}
fn parse_klipy_path(raw_url: &str) -> Option<KlipyPath> {
let parsed = Url::parse(raw_url).ok()?;
let hostname = parsed.host_str()?.to_ascii_lowercase();
if hostname != "klipy.com" && hostname != "www.klipy.com" {
return None;
}
let mut segments = parsed.path_segments()?;
let kind = segments.next()?.to_ascii_lowercase();
let slug = segments.next()?.trim().to_owned();
if slug.is_empty() {
return None;
}
let path_type = match kind.as_str() {
"gif" | "gifs" => KlipyPathType::Gif,
"clip" | "clips" => KlipyPathType::Clip,
_ => return None,
};
Some(KlipyPath { path_type, slug })
}
fn build_share_url_with_type(path_type: KlipyPathType, slug: &str) -> String {
let base_path = match path_type {
KlipyPathType::Gif => "gifs",
KlipyPathType::Clip => "clips",
};
let encoded_slug = urlencoding::encode(slug);
format!("https://klipy.com/{base_path}/{encoded_slug}")
}
fn klipy_resource(path_type: KlipyPathType) -> &'static str {
match path_type {
KlipyPathType::Gif => "gifs",
KlipyPathType::Clip => "clips",
}
}
fn public_format_key(size: &str, format: &str) -> String {
match (size, format) {
("hd", "webm") => "webm",
("hd", "mp4") => "mp4",
("hd", "webp") => "webp",
("hd", "gif") => "gif",
("md", "webm") => "mediumwebm",
("md", "mp4") => "mediummp4",
("md", "webp") => "mediumwebp",
("md", "gif") => "mediumgif",
("sm", "webm") => "tinywebm",
("sm", "mp4") => "tinymp4",
("sm", "webp") => "tinywebp",
("sm", "gif") => "tinygif",
("xs", "webm") => "nanowebm",
("xs", "mp4") => "nanomp4",
("xs", "webp") => "nanowebp",
("xs", "gif") => "nanogif",
_ => format,
}
.to_owned()
}
fn category_response(name: String, gif: Option<GifItem>) -> GifCategoryTag {
GifCategoryTag {
src: gif.as_ref().map(|gif| gif.src.clone()).unwrap_or_default(),
proxy_src: gif
.as_ref()
.map(|gif| gif.proxy_src.clone())
.unwrap_or_default(),
gif,
name,
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn locale_uses_klipy_underscore_form() {
assert_eq!(normalize_locale("en-US"), "en_US");
assert_eq!(normalize_locale("sv_SE"), "sv_SE");
}
#[test]
fn extracts_klipy_slug_only_from_klipy_hosts() {
assert_eq!(
extract_slug_from_url("https://klipy.com/gifs/funny-123").as_deref(),
Some("funny-123")
);
assert_eq!(
extract_slug_from_url("https://www.klipy.com/clip/abc").as_deref(),
Some("abc")
);
assert_eq!(
extract_slug_from_url("https://notklipy.com/gifs/funny"),
None
);
}
#[test]
fn build_share_url_uses_gifs_path() {
assert_eq!(build_share_url("hello"), "https://klipy.com/gifs/hello");
assert_eq!(build_share_url(" "), "https://klipy.com/gifs");
assert_eq!(
build_share_url_with_type(KlipyPathType::Clip, "hello"),
"https://klipy.com/clips/hello"
);
}
#[test]
fn stringifies_numeric_klipy_ids() {
assert_eq!(
klipy_id_as_string(&serde_json::json!(2484942301552561_i64)).as_deref(),
Some("2484942301552561")
);
assert_eq!(
klipy_id_as_string(&serde_json::json!(" abc ")).as_deref(),
Some("abc")
);
assert_eq!(klipy_id_as_string(&serde_json::json!(" ")), None);
}
#[test]
fn maps_provider_format_keys() {
assert_eq!(public_format_key("hd", "webm"), "webm");
assert_eq!(public_format_key("sm", "gif"), "tinygif");
assert_eq!(public_format_key("xs", "webp"), "nanowebp");
}
#[tokio::test]
#[ignore]
async fn live_resolves_klipy_url_with_direct_lookup() {
let api_key = std::env::var("FLUXER_KLIPY_API_KEY")
.or_else(|_| std::env::var("KLIPY_API_KEY"))
.expect("FLUXER_KLIPY_API_KEY or KLIPY_API_KEY set");
let client =
KlipyClient::new(MediaProxyUrlBuilder::from_env().expect("media proxy env configured"))
.expect("KLIPY client");
let gif = client
.resolve_by_url(
&api_key,
"https://klipy.com/gifs/goatplaybanjo-chat-4",
"en-US",
"US",
)
.await
.expect("KLIPY direct lookup")
.expect("resolved GIF");
assert_eq!(gif.slug, "goatplaybanjo-chat-4");
assert_eq!(gif.provider, KLIPY_PROVIDER_NAME);
assert_eq!(gif.url, "https://klipy.com/gifs/goatplaybanjo-chat-4");
assert!(gif.width > 0);
assert!(gif.height > 0);
assert!(gif.media.contains_key("webm") || gif.media.contains_key("mp4"));
assert!(gif.proxy_src.starts_with("http"));
}
}
+39
View File
@@ -0,0 +1,39 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
mod klipy;
mod media_proxy;
mod router_impl;
mod shard_impl;
mod types;
use fluxer_svc::config::{Mode, ServiceConfig};
use fluxer_svc::transport::NatsTransport;
use router_impl::GifsRouter;
use shard_impl::GifsShard;
#[tokio::main]
async fn main() -> anyhow::Result<()> {
fluxer_svc::init_tracing();
let config = ServiceConfig::from_env()?;
let transport = NatsTransport::connect(&config.nats_url).await?;
tracing::info!(
service = config.service_name,
mode = ?config.mode,
shard_id = config.shard_id,
shard_count = config.shard_count,
listen_addr = %config.listen_addr,
"starting gifs service"
);
match config.mode {
Mode::Router => {
let router = GifsRouter::new(config.cache_max_entries, config.cache_ttl);
fluxer_svc::router::run_router(&config, router, transport).await
}
Mode::Shard => {
let shard = GifsShard::new(&config)?;
fluxer_svc::shard::run_shard(&config, shard, transport).await
}
}
}
+216
View File
@@ -0,0 +1,216 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
use base64::prelude::*;
use hmac::{Hmac, KeyInit, Mac};
use sha2::Sha256;
use url::Url;
const V2_PATH_PREFIX: &str = "v2/";
#[derive(Clone)]
pub struct MediaProxyUrlBuilder {
endpoint: String,
endpoint_host: Option<String>,
secret_key: String,
}
impl MediaProxyUrlBuilder {
pub fn from_env() -> anyhow::Result<Self> {
let endpoint = std::env::var("FLUXER_MEDIA_PROXY_PUBLIC_ENDPOINT")
.or_else(|_| std::env::var("FLUXER_MEDIA_ENDPOINT"))
.unwrap_or_default();
if endpoint.trim().is_empty() {
anyhow::bail!(
"gifs shard requires FLUXER_MEDIA_PROXY_PUBLIC_ENDPOINT or FLUXER_MEDIA_ENDPOINT"
);
}
let secret_key = std::env::var("FLUXER_MEDIA_PROXY_SECRET_KEY").unwrap_or_default();
if secret_key.trim().is_empty() {
anyhow::bail!("gifs shard requires FLUXER_MEDIA_PROXY_SECRET_KEY");
}
let endpoint = endpoint.trim_end_matches('/').to_owned();
let endpoint_host = Url::parse(&endpoint)
.ok()
.and_then(|parsed| parsed.host_str().map(ToOwned::to_owned));
Ok(Self {
endpoint,
endpoint_host,
secret_key,
})
}
pub fn external_proxy_url(&self, input_url: &str) -> Option<String> {
let parsed = Url::parse(input_url).ok()?;
if self
.endpoint_host
.as_deref()
.is_some_and(|host| parsed.host_str() == Some(host))
{
return Some(input_url.to_owned());
}
let proxy_path = build_external_media_proxy_path(parsed.as_str());
let signature = create_signature(&proxy_path, &self.secret_key);
Some(format!(
"{}/external/{signature}/{proxy_path}",
self.endpoint
))
}
}
fn build_external_media_proxy_path(input_url: &str) -> String {
format!(
"{V2_PATH_PREFIX}{}",
BASE64_URL_SAFE_NO_PAD.encode(input_url)
)
}
fn create_signature(input: &str, secret: &str) -> String {
let mut mac = Hmac::<Sha256>::new_from_slice(secret.as_bytes()).expect("HMAC accepts any key");
mac.update(input.as_bytes());
BASE64_URL_SAFE_NO_PAD.encode(mac.finalize().into_bytes())
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Mutex;
static ENV_LOCK: Mutex<()> = Mutex::new(());
fn with_media_proxy_env(
vars: &[(&str, Option<&str>)],
test: impl FnOnce() -> anyhow::Result<()>,
) -> anyhow::Result<()> {
let _guard = ENV_LOCK.lock().unwrap();
let keys = [
"FLUXER_MEDIA_PROXY_PUBLIC_ENDPOINT",
"FLUXER_MEDIA_ENDPOINT",
"FLUXER_MEDIA_PROXY_ENDPOINT",
"FLUXER_MEDIA_PROXY_SECRET_KEY",
];
let saved = keys
.iter()
.map(|key| (*key, std::env::var(key).ok()))
.collect::<Vec<_>>();
for key in keys {
unsafe {
std::env::remove_var(key);
}
}
for (key, value) in vars {
if let Some(value) = value {
unsafe {
std::env::set_var(key, value);
}
}
}
let result = test();
for (key, value) in saved {
match value {
Some(value) => unsafe {
std::env::set_var(key, value);
},
None => unsafe {
std::env::remove_var(key);
},
}
}
result
}
#[test]
fn external_proxy_url_builds_v2_signed_url() {
let builder = MediaProxyUrlBuilder {
endpoint: "https://media.example.test".to_owned(),
endpoint_host: Some("media.example.test".to_owned()),
secret_key: "secret".to_owned(),
};
let url = builder
.external_proxy_url("https://img.klipy.com/a.webp?x=1")
.expect("proxy url");
assert!(url.starts_with("https://media.example.test/external/"));
assert!(url.contains("/v2/"));
assert_eq!(
builder.external_proxy_url("https://media.example.test/external/existing"),
Some("https://media.example.test/external/existing".to_owned())
);
}
#[test]
fn from_env_uses_public_endpoint_when_internal_proxy_endpoint_is_set() -> anyhow::Result<()> {
with_media_proxy_env(
&[
(
"FLUXER_MEDIA_PROXY_PUBLIC_ENDPOINT",
Some("https://media.example.test/"),
),
(
"FLUXER_MEDIA_PROXY_ENDPOINT",
Some("http://media-proxy:8080"),
),
("FLUXER_MEDIA_PROXY_SECRET_KEY", Some("secret")),
],
|| {
let builder = MediaProxyUrlBuilder::from_env()?;
assert_eq!(builder.endpoint, "https://media.example.test");
Ok(())
},
)
}
#[test]
fn from_env_accepts_legacy_public_media_endpoint() -> anyhow::Result<()> {
with_media_proxy_env(
&[
(
"FLUXER_MEDIA_ENDPOINT",
Some("https://media.example.test/media"),
),
(
"FLUXER_MEDIA_PROXY_ENDPOINT",
Some("http://media-proxy:8080"),
),
("FLUXER_MEDIA_PROXY_SECRET_KEY", Some("secret")),
],
|| {
let builder = MediaProxyUrlBuilder::from_env()?;
assert_eq!(builder.endpoint, "https://media.example.test/media");
Ok(())
},
)
}
#[test]
fn from_env_rejects_internal_proxy_endpoint_without_public_endpoint() -> anyhow::Result<()> {
with_media_proxy_env(
&[
(
"FLUXER_MEDIA_PROXY_ENDPOINT",
Some("http://media-proxy:8080"),
),
("FLUXER_MEDIA_PROXY_SECRET_KEY", Some("secret")),
],
|| {
let err = MediaProxyUrlBuilder::from_env()
.err()
.expect("internal endpoint must not be accepted as public endpoint")
.to_string();
assert!(err.contains("FLUXER_MEDIA_PROXY_PUBLIC_ENDPOINT"));
Ok(())
},
)
}
}
+183
View File
@@ -0,0 +1,183 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
use crate::types::{GifRequest, GifServiceResponse};
use fluxer_svc::router::RouterService;
use moka::sync::Cache;
use std::time::Duration;
pub struct GifsRouter {
l1: Cache<String, GifServiceResponse>,
}
impl GifsRouter {
pub fn new(max_entries: u64, ttl: Duration) -> Self {
Self {
l1: Cache::builder()
.max_capacity(max_entries)
.time_to_live(ttl)
.build(),
}
}
}
impl RouterService for GifsRouter {
type Request = GifRequest;
type Response = GifServiceResponse;
fn service_name(&self) -> &str {
"gifs"
}
fn route_key(req: &GifRequest) -> String {
request_key(req).unwrap_or_else(|| "uncached".to_owned())
}
fn coalesce_key(req: &GifRequest) -> Option<String> {
request_key(req)
}
fn l1_lookup(&self, req: &GifRequest) -> Option<GifServiceResponse> {
request_l1_key(req).and_then(|key| self.l1.get(&key))
}
fn l1_insert(&self, req: &GifRequest, resp: &GifServiceResponse) {
if response_is_cacheable(resp)
&& let Some(key) = request_l1_key(req)
{
self.l1.insert(key, resp.clone());
}
}
fn l1_invalidate(&self, key: &str) {
self.l1.invalidate(key);
}
}
fn request_l1_key(req: &GifRequest) -> Option<String> {
match req {
GifRequest::Search { .. }
| GifRequest::GetFeatured { .. }
| GifRequest::GetTrendingGifs { .. }
| GifRequest::Suggest { .. }
| GifRequest::ResolveByUrl { .. }
| GifRequest::BuildShareUrl { .. }
| GifRequest::ExtractSlugFromUrl { .. } => request_key(req),
GifRequest::IsAvailable { .. } | GifRequest::RegisterShare { .. } => None,
}
}
fn request_key(req: &GifRequest) -> Option<String> {
match req {
GifRequest::IsAvailable { .. } => None,
GifRequest::Search {
q, locale, country, ..
} => Some(format!("search:{locale}:{country}:{q}")),
GifRequest::GetFeatured {
locale, country, ..
} => Some(format!("featured:{locale}:{country}")),
GifRequest::GetTrendingGifs {
locale, country, ..
} => Some(format!("trending:{locale}:{country}")),
GifRequest::Suggest { q, locale, .. } => Some(format!("suggest:{locale}:{q}")),
GifRequest::RegisterShare { .. } => None,
GifRequest::ResolveByUrl {
url,
locale,
country,
..
} => Some(format!("resolve:{locale}:{country}:{url}")),
GifRequest::BuildShareUrl { slug } => Some(format!("share-url:{slug}")),
GifRequest::ExtractSlugFromUrl { url } => Some(format!("extract-slug:{url}")),
}
}
fn response_is_cacheable(resp: &GifServiceResponse) -> bool {
!matches!(
resp,
GifServiceResponse::Available { .. }
| GifServiceResponse::Registered
| GifServiceResponse::Failed { .. }
)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn coalesce_key_excludes_api_key() {
let first = GifRequest::Search {
api_key: "a".to_owned(),
q: "wave".to_owned(),
locale: "en_US".to_owned(),
country: "US".to_owned(),
};
let second = GifRequest::Search {
api_key: "b".to_owned(),
q: "wave".to_owned(),
locale: "en_US".to_owned(),
country: "US".to_owned(),
};
assert_eq!(
GifsRouter::coalesce_key(&first),
GifsRouter::coalesce_key(&second)
);
}
#[test]
fn trending_key_is_locale_country_scoped_and_api_key_free() {
let base = GifRequest::GetTrendingGifs {
api_key: "a".to_owned(),
locale: "en_US".to_owned(),
country: "US".to_owned(),
};
let different_api_key = GifRequest::GetTrendingGifs {
api_key: "b".to_owned(),
locale: "en_US".to_owned(),
country: "US".to_owned(),
};
let different_locale = GifRequest::GetTrendingGifs {
api_key: "a".to_owned(),
locale: "sv_SE".to_owned(),
country: "US".to_owned(),
};
let different_country = GifRequest::GetTrendingGifs {
api_key: "a".to_owned(),
locale: "en_US".to_owned(),
country: "SE".to_owned(),
};
assert_eq!(
GifsRouter::coalesce_key(&base),
Some("trending:en_US:US".to_owned())
);
assert_eq!(
GifsRouter::coalesce_key(&base),
GifsRouter::coalesce_key(&different_api_key)
);
assert_ne!(
GifsRouter::coalesce_key(&base),
GifsRouter::coalesce_key(&different_locale)
);
assert_ne!(
GifsRouter::coalesce_key(&base),
GifsRouter::coalesce_key(&different_country)
);
assert_eq!(request_l1_key(&base), Some("trending:en_US:US".to_owned()));
}
#[test]
fn register_share_is_not_cached_or_coalesced() {
let request = GifRequest::RegisterShare {
api_key: "key".to_owned(),
id: "gif".to_owned(),
q: "wave".to_owned(),
locale: "en_US".to_owned(),
country: "US".to_owned(),
};
assert_eq!(GifsRouter::coalesce_key(&request), None);
assert_eq!(request_l1_key(&request), None);
}
}
+410
View File
@@ -0,0 +1,410 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
use crate::klipy::{KlipyClient, build_share_url, extract_slug_from_url};
use crate::media_proxy::MediaProxyUrlBuilder;
use crate::types::{GifCategoryTag, GifItem, GifRequest, GifServiceResponse};
use fluxer_svc::config::ServiceConfig;
use fluxer_svc::shard::ShardService;
use moka::future::Cache;
use std::collections::HashSet;
use std::future::Future;
use std::sync::Arc;
use std::time::{Duration, Instant};
use tokio::sync::Mutex;
const SEARCH_SOFT_TTL: Duration = Duration::from_secs(30);
const SEARCH_HARD_TTL: Duration = Duration::from_secs(5 * 60);
const SUGGEST_SOFT_TTL: Duration = Duration::from_secs(60);
const SUGGEST_HARD_TTL: Duration = Duration::from_secs(10 * 60);
const FEATURED_GIFS_SOFT_TTL: Duration = Duration::from_secs(5 * 60);
const FEATURED_GIFS_HARD_TTL: Duration = Duration::from_secs(30 * 60);
const CATEGORIES_SOFT_TTL: Duration = Duration::from_secs(24 * 60 * 60);
const CATEGORIES_HARD_TTL: Duration = Duration::from_secs(48 * 60 * 60);
const RESOLVE_SOFT_TTL: Duration = Duration::from_secs(30 * 60);
const RESOLVE_HARD_TTL: Duration = Duration::from_secs(2 * 60 * 60);
#[derive(Clone)]
pub struct GifsShard {
inner: Arc<GifsShardInner>,
}
struct GifsShardInner {
klipy: KlipyClient,
gif_lists: Cache<String, Cached<Vec<GifItem>>>,
categories: Cache<String, Cached<Vec<GifCategoryTag>>>,
suggestions: Cache<String, Cached<Vec<String>>>,
resolved: Cache<String, Cached<Option<GifItem>>>,
refreshing: Mutex<HashSet<String>>,
}
#[derive(Debug, Clone)]
struct Cached<T> {
data: T,
stored_at: Instant,
}
#[derive(Debug, Clone, Copy)]
struct CachePolicy {
soft_ttl: Duration,
hard_ttl: Duration,
}
impl CachePolicy {
const fn new(soft_ttl: Duration, hard_ttl: Duration) -> Self {
Self { soft_ttl, hard_ttl }
}
}
impl<T> Cached<T> {
fn new(data: T) -> Self {
Self {
data,
stored_at: Instant::now(),
}
}
fn age(&self) -> Duration {
self.stored_at.elapsed()
}
}
impl GifsShard {
pub fn new(config: &ServiceConfig) -> anyhow::Result<Self> {
let media_proxy = MediaProxyUrlBuilder::from_env()?;
let klipy = KlipyClient::new(media_proxy)?;
let max_capacity = config.cache_max_entries;
let max_cache_ttl = CATEGORIES_HARD_TTL;
Ok(Self {
inner: Arc::new(GifsShardInner {
klipy,
gif_lists: build_cache(max_capacity, max_cache_ttl),
categories: build_cache(max_capacity, max_cache_ttl),
suggestions: build_cache(max_capacity, max_cache_ttl),
resolved: build_cache(max_capacity, max_cache_ttl),
refreshing: Mutex::new(HashSet::new()),
}),
})
}
async fn get_cached<T, Fetch, Fut>(
&self,
cache: Cache<String, Cached<T>>,
key: String,
policy: CachePolicy,
fetch: Fetch,
) -> anyhow::Result<T>
where
T: Clone + Send + Sync + 'static,
Fetch: Fn() -> Fut + Clone + Send + Sync + 'static,
Fut: Future<Output = anyhow::Result<T>> + Send + 'static,
{
if let Some(cached) = cache.get(&key).await {
let age = cached.age();
if age <= policy.soft_ttl {
return Ok(cached.data);
}
if age <= policy.hard_ttl {
self.trigger_background_refresh(cache, key, policy, fetch);
return Ok(cached.data);
}
cache.invalidate(&key).await;
}
let fetch_for_load = fetch.clone();
let cached = cache
.try_get_with(key, async move {
let data = fetch_for_load().await?;
Ok::<Cached<T>, anyhow::Error>(Cached::new(data))
})
.await
.map_err(|error| anyhow::anyhow!("{}", error.as_ref()))?;
Ok(cached.data)
}
fn trigger_background_refresh<T, Fetch, Fut>(
&self,
cache: Cache<String, Cached<T>>,
key: String,
_policy: CachePolicy,
fetch: Fetch,
) where
T: Clone + Send + Sync + 'static,
Fetch: Fn() -> Fut + Clone + Send + Sync + 'static,
Fut: Future<Output = anyhow::Result<T>> + Send + 'static,
{
let this = self.clone();
tokio::spawn(async move {
{
let mut refreshing = this.inner.refreshing.lock().await;
if !refreshing.insert(key.clone()) {
return;
}
}
let result = fetch().await;
match result {
Ok(data) => {
cache.insert(key.clone(), Cached::new(data)).await;
}
Err(error) => {
tracing::debug!(error = %error, cache_key = %key, "background GIF cache refresh failed");
}
}
let mut refreshing = this.inner.refreshing.lock().await;
refreshing.remove(&key);
});
}
async fn handle_available(&self, api_key: Option<String>) -> GifServiceResponse {
GifServiceResponse::Available {
available: api_key.as_deref().is_some_and(|key| !key.trim().is_empty()),
}
}
async fn handle_search(
&self,
api_key: String,
q: String,
locale: String,
country: String,
) -> anyhow::Result<GifServiceResponse> {
let key = format!("search:{locale}:{country}:{q}");
let this = self.clone();
let gifs = self
.get_cached(
self.inner.gif_lists.clone(),
key,
CachePolicy::new(SEARCH_SOFT_TTL, SEARCH_HARD_TTL),
move || {
let this = this.clone();
let api_key = api_key.clone();
let q = q.clone();
let locale = locale.clone();
let country = country.clone();
async move {
this.inner
.klipy
.search(&api_key, &q, &locale, &country, 50)
.await
}
},
)
.await?;
Ok(GifServiceResponse::SearchResults(gifs))
}
async fn handle_featured(
&self,
api_key: String,
locale: String,
country: String,
) -> anyhow::Result<GifServiceResponse> {
let gifs_key = format!("featured_gifs:{locale}:{country}");
let categories_key = format!("featured_categories:{locale}");
let gifs_this = self.clone();
let categories_this = self.clone();
let api_key_for_gifs = api_key.clone();
let locale_for_gifs = locale.clone();
let country_for_gifs = country.clone();
let api_key_for_categories = api_key;
let locale_for_categories = locale;
let gifs_future = self.get_cached(
self.inner.gif_lists.clone(),
gifs_key,
CachePolicy::new(FEATURED_GIFS_SOFT_TTL, FEATURED_GIFS_HARD_TTL),
move || {
let this = gifs_this.clone();
let api_key = api_key_for_gifs.clone();
let locale = locale_for_gifs.clone();
let country = country_for_gifs.clone();
async move {
this.inner
.klipy
.featured_gifs(&api_key, &locale, &country)
.await
}
},
);
let categories_future = self.get_cached(
self.inner.categories.clone(),
categories_key,
CachePolicy::new(CATEGORIES_SOFT_TTL, CATEGORIES_HARD_TTL),
move || {
let this = categories_this.clone();
let api_key = api_key_for_categories.clone();
let locale = locale_for_categories.clone();
async move {
this.inner
.klipy
.featured_categories(&api_key, &locale)
.await
}
},
);
let (gifs, categories) = tokio::try_join!(gifs_future, categories_future)?;
Ok(GifServiceResponse::Featured { gifs, categories })
}
async fn handle_trending(
&self,
api_key: String,
locale: String,
country: String,
) -> anyhow::Result<GifServiceResponse> {
let key = format!("trending:{locale}:{country}");
let this = self.clone();
let gifs = self
.get_cached(
self.inner.gif_lists.clone(),
key,
CachePolicy::new(FEATURED_GIFS_SOFT_TTL, FEATURED_GIFS_HARD_TTL),
move || {
let this = this.clone();
let api_key = api_key.clone();
let locale = locale.clone();
let country = country.clone();
async move {
this.inner
.klipy
.trending_gifs(&api_key, &locale, &country)
.await
}
},
)
.await?;
Ok(GifServiceResponse::TrendingResults(gifs))
}
async fn handle_suggest(
&self,
api_key: String,
q: String,
locale: String,
) -> anyhow::Result<GifServiceResponse> {
let key = format!("suggest:{locale}:{q}");
let this = self.clone();
let suggestions = self
.get_cached(
self.inner.suggestions.clone(),
key,
CachePolicy::new(SUGGEST_SOFT_TTL, SUGGEST_HARD_TTL),
move || {
let this = this.clone();
let api_key = api_key.clone();
let q = q.clone();
let locale = locale.clone();
async move { this.inner.klipy.suggestions(&api_key, &q, &locale).await }
},
)
.await?;
Ok(GifServiceResponse::Suggestions(suggestions))
}
async fn handle_resolve_by_url(
&self,
api_key: String,
url: String,
locale: String,
country: String,
) -> anyhow::Result<GifServiceResponse> {
let key = format!("resolve:{locale}:{country}:{url}");
let this = self.clone();
let gif = self
.get_cached(
self.inner.resolved.clone(),
key,
CachePolicy::new(RESOLVE_SOFT_TTL, RESOLVE_HARD_TTL),
move || {
let this = this.clone();
let api_key = api_key.clone();
let url = url.clone();
let locale = locale.clone();
let country = country.clone();
async move {
this.inner
.klipy
.resolve_by_url(&api_key, &url, &locale, &country)
.await
}
},
)
.await?;
Ok(GifServiceResponse::Resolved { gif })
}
}
impl ShardService for GifsShard {
type Request = GifRequest;
type Response = GifServiceResponse;
fn service_name(&self) -> &str {
"gifs"
}
async fn handle(&self, request: GifRequest) -> anyhow::Result<GifServiceResponse> {
let response = match request {
GifRequest::IsAvailable { api_key } => self.handle_available(api_key).await,
GifRequest::Search {
api_key,
q,
locale,
country,
} => self.handle_search(api_key, q, locale, country).await?,
GifRequest::GetFeatured {
api_key,
locale,
country,
} => self.handle_featured(api_key, locale, country).await?,
GifRequest::GetTrendingGifs {
api_key,
locale,
country,
} => self.handle_trending(api_key, locale, country).await?,
GifRequest::Suggest { api_key, q, locale } => {
self.handle_suggest(api_key, q, locale).await?
}
GifRequest::RegisterShare {
api_key,
id,
q,
locale,
country,
} => {
self.inner
.klipy
.register_share(&api_key, &id, &q, &locale, &country)
.await?;
GifServiceResponse::Registered
}
GifRequest::ResolveByUrl {
api_key,
url,
locale,
country,
} => {
self.handle_resolve_by_url(api_key, url, locale, country)
.await?
}
GifRequest::BuildShareUrl { slug } => GifServiceResponse::ShareUrl {
url: build_share_url(&slug),
},
GifRequest::ExtractSlugFromUrl { url } => GifServiceResponse::ExtractedSlug {
slug: extract_slug_from_url(&url),
},
};
Ok(response)
}
}
fn build_cache<T>(max_capacity: u64, time_to_live: Duration) -> Cache<String, Cached<T>>
where
T: Clone + Send + Sync + 'static,
{
Cache::builder()
.max_capacity(max_capacity)
.time_to_live(time_to_live)
.build()
}
+111
View File
@@ -0,0 +1,111 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
use serde::{Deserialize, Serialize};
use std::collections::BTreeMap;
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "op", rename_all_fields = "snake_case")]
pub enum GifRequest {
IsAvailable {
api_key: Option<String>,
},
Search {
api_key: String,
q: String,
locale: String,
country: String,
},
GetFeatured {
api_key: String,
locale: String,
country: String,
},
GetTrendingGifs {
api_key: String,
locale: String,
country: String,
},
Suggest {
api_key: String,
q: String,
locale: String,
},
RegisterShare {
api_key: String,
id: String,
q: String,
locale: String,
country: String,
},
ResolveByUrl {
api_key: String,
url: String,
locale: String,
country: String,
},
BuildShareUrl {
slug: String,
},
ExtractSlugFromUrl {
url: String,
},
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub enum GifServiceResponse {
Available {
available: bool,
},
SearchResults(Vec<GifItem>),
Featured {
gifs: Vec<GifItem>,
categories: Vec<GifCategoryTag>,
},
TrendingResults(Vec<GifItem>),
Suggestions(Vec<String>),
Registered,
Resolved {
gif: Option<GifItem>,
},
ShareUrl {
url: String,
},
ExtractedSlug {
slug: Option<String>,
},
Failed {
message: String,
},
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct GifMediaFormat {
pub src: String,
pub proxy_src: String,
pub width: i32,
pub height: i32,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct GifItem {
pub id: String,
pub slug: String,
pub provider: String,
pub title: String,
pub url: String,
pub src: String,
pub proxy_src: String,
pub width: i32,
pub height: i32,
pub media: BTreeMap<String, GifMediaFormat>,
#[serde(skip_serializing_if = "Option::is_none")]
pub placeholder: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct GifCategoryTag {
pub name: String,
pub src: String,
pub proxy_src: String,
pub gif: Option<GifItem>,
}