From 5c51ec895f5f697f6d32a23fc282fd4564dac5eb Mon Sep 17 00:00:00 2001 From: asepharyana Date: Wed, 22 Jul 2026 14:51:08 +0700 Subject: [PATCH] feat: remove croxy proxy and image cache features --- Cargo.lock | 210 +-- Cargo.toml | 5 - src/application/anime/use_cases.rs | 75 +- src/application/anime2/use_cases.rs | 139 +- src/application/komik/use_cases.rs | 89 +- src/application/proxy/use_cases.rs | 235 ---- src/bootstrap/mod.rs | 11 - src/bootstrap/setup.rs | 44 +- src/config/mod.rs | 22 - src/domain/entity/anime.rs | 111 -- src/domain/entity/komik.rs | 11 - src/domain/repository/image_cache.rs | 19 - src/domain/repository/mod.rs | 2 - src/events/bus.rs | 9 - src/infrastructure/mod.rs | 2 - .../persistence/entities/image_cache.rs | 71 - .../persistence/entities/mod.rs | 1 - src/infrastructure/persistence/mod.rs | 1 - .../repository/image_cache_seaorm.rs | 131 -- src/infrastructure/repository/mod.rs | 2 - .../repository/parsers/alqanime_parser.rs | 15 +- src/infrastructure/services/images/cache.rs | 1142 ----------------- src/infrastructure/services/images/mod.rs | 1 - src/infrastructure/services/mod.rs | 1 - src/lib.rs | 1 - src/observability/openapi_modules.rs | 2 - src/presentation/handler/anime.rs | 7 +- src/presentation/handler/anime2.rs | 7 +- src/presentation/handler/komik.rs | 7 +- src/presentation/handler/proxy.rs | 132 +- src/presentation/router.rs | 13 +- src/presentation/state.rs | 3 - src/scheduler/cleanup_cache.rs | 222 ---- src/scheduler/mod.rs | 5 - src/scheduler/runner.rs | 93 -- 35 files changed, 42 insertions(+), 2799 deletions(-) delete mode 100644 src/domain/repository/image_cache.rs delete mode 100644 src/infrastructure/persistence/entities/image_cache.rs delete mode 100644 src/infrastructure/persistence/entities/mod.rs delete mode 100644 src/infrastructure/persistence/mod.rs delete mode 100644 src/infrastructure/repository/image_cache_seaorm.rs delete mode 100644 src/infrastructure/services/images/cache.rs delete mode 100644 src/infrastructure/services/images/mod.rs delete mode 100644 src/infrastructure/services/mod.rs delete mode 100644 src/scheduler/cleanup_cache.rs delete mode 100644 src/scheduler/mod.rs delete mode 100644 src/scheduler/runner.rs diff --git a/Cargo.lock b/Cargo.lock index 03a9829..61ef137 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -439,17 +439,6 @@ dependencies = [ "shlex", ] -[[package]] -name = "cfb" -version = "0.7.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d38f2da7a0a2c4ccf0065be06397cc26a81f4e528be095826eee9d4adbb8c60f" -dependencies = [ - "byteorder", - "fnv", - "uuid", -] - [[package]] name = "cfg-if" version = "1.0.4" @@ -476,16 +465,6 @@ dependencies = [ "windows-link", ] -[[package]] -name = "chrono-tz" -version = "0.10.4" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a6139a8597ed92cf816dfb33f5dd6cf0bb93a6adc938f11039f371bc5bcd26c3" -dependencies = [ - "chrono", - "phf 0.12.1", -] - [[package]] name = "combine" version = "4.6.7" @@ -643,17 +622,6 @@ dependencies = [ "cfg-if", ] -[[package]] -name = "croner" -version = "3.0.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4aa42bcd3d846ebf66e15bd528d1087f75d1c6c1c66ebff626178a106353c576" -dependencies = [ - "chrono", - "derive_builder", - "strum 0.27.2", -] - [[package]] name = "crossbeam-deque" version = "0.8.6" @@ -713,7 +681,7 @@ dependencies = [ "cssparser-macros", "dtoa-short", "itoa", - "phf 0.13.1", + "phf", "smallvec", ] @@ -727,41 +695,6 @@ dependencies = [ "syn 2.0.117", ] -[[package]] -name = "darling" -version = "0.20.11" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "fc7f46116c46ff9ab3eb1597a45688b6715c6e628b5c133e288e709a29bcb4ee" -dependencies = [ - "darling_core", - "darling_macro", -] - -[[package]] -name = "darling_core" -version = "0.20.11" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0d00b9596d185e565c2207a0b01f8bd1a135483d02d9b7b0a54b11da8d53412e" -dependencies = [ - "fnv", - "ident_case", - "proc-macro2", - "quote", - "strsim", - "syn 2.0.117", -] - -[[package]] -name = "darling_macro" -version = "0.20.11" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "fc34b93ccb385b40dc71c6fceac4b2ad23662c7eeb248cf10d529b7e055b6ead" -dependencies = [ - "darling_core", - "quote", - "syn 2.0.117", -] - [[package]] name = "dashmap" version = "6.1.0" @@ -847,37 +780,6 @@ dependencies = [ "syn 2.0.117", ] -[[package]] -name = "derive_builder" -version = "0.20.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "507dfb09ea8b7fa618fcf76e953f4f5e192547945816d5358edffe39f6f94947" -dependencies = [ - "derive_builder_macro", -] - -[[package]] -name = "derive_builder_core" -version = "0.20.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2d5bcf7b024d6835cfb3d473887cd966994907effbe9227e8c8219824d06c4e8" -dependencies = [ - "darling", - "proc-macro2", - "quote", - "syn 2.0.117", -] - -[[package]] -name = "derive_builder_macro" -version = "0.20.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "ab63b0e2bf4d5928aff72e83a7dace85d7bba5fe12dcc3c5a572d78caffd3f3c" -dependencies = [ - "derive_builder_core", - "syn 2.0.117", -] - [[package]] name = "derive_more" version = "2.1.1" @@ -1648,12 +1550,6 @@ version = "2.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "3d3067d79b975e8844ca9eb072e16b31c3c1c36928edf9c6789548c524d0d954" -[[package]] -name = "ident_case" -version = "1.0.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b9e0384b61958566e926dc50660321d12159025e767c18e043daf26b70104c39" - [[package]] name = "idna" version = "1.1.0" @@ -1697,15 +1593,6 @@ dependencies = [ "serde_core", ] -[[package]] -name = "infer" -version = "0.19.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a588916bfdfd92e71cacef98a63d9b1f0d74d6599980d11894290e7ddefffcf7" -dependencies = [ - "cfb", -] - [[package]] name = "inherent" version = "1.0.13" @@ -2038,17 +1925,6 @@ version = "0.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "51d515d32fb182ee37cda2ccdcb92950d6a3c2893aa280e540671c2cd0f3b1d9" -[[package]] -name = "num-derive" -version = "0.4.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "ed3955f1a9c7c0c15e092f9c887db08b1fc683305fdf6eb6684f22555355e202" -dependencies = [ - "proc-macro2", - "quote", - "syn 2.0.117", -] - [[package]] name = "num-integer" version = "0.1.46" @@ -2350,15 +2226,6 @@ dependencies = [ "serde", ] -[[package]] -name = "phf" -version = "0.12.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "913273894cec178f401a31ec4b656318d95473527be05c0752cc41cdc32be8b7" -dependencies = [ - "phf_shared 0.12.1", -] - [[package]] name = "phf" version = "0.13.1" @@ -2366,7 +2233,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c1562dc717473dbaa4c1f85a36410e03c047b2e7df7f45ee938fbef64ae7fadf" dependencies = [ "phf_macros", - "phf_shared 0.13.1", + "phf_shared", "serde", ] @@ -2377,7 +2244,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "49aa7f9d80421bca176ca8dbfebe668cc7a2684708594ec9f3c0db0805d5d6e1" dependencies = [ "phf_generator", - "phf_shared 0.13.1", + "phf_shared", ] [[package]] @@ -2387,7 +2254,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "135ace3a761e564ec88c03a77317a7c6b80bb7f7135ef2544dbe054243b89737" dependencies = [ "fastrand", - "phf_shared 0.13.1", + "phf_shared", ] [[package]] @@ -2397,21 +2264,12 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "812f032b54b1e759ccd5f8b6677695d5268c588701effba24601f6932f8269ef" dependencies = [ "phf_generator", - "phf_shared 0.13.1", + "phf_shared", "proc-macro2", "quote", "syn 2.0.117", ] -[[package]] -name = "phf_shared" -version = "0.12.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "06005508882fb681fd97892ecff4b7fd0fee13ef1aa569f8695dae7ab9099981" -dependencies = [ - "siphasher", -] - [[package]] name = "phf_shared" version = "0.13.1" @@ -3124,7 +2982,6 @@ dependencies = [ "async-trait", "axum 0.8.8", "backoff", - "bytes", "chrono", "config", "dashmap", @@ -3132,9 +2989,7 @@ dependencies = [ "dotenvy", "flate2", "futures", - "hex", "http", - "infer", "opentelemetry", "opentelemetry-otlp", "opentelemetry_sdk", @@ -3146,10 +3001,8 @@ dependencies = [ "sea-orm", "serde", "serde_json", - "sha2", "thiserror 2.0.18", "tokio", - "tokio-cron-scheduler", "tower-http", "tracing", "tracing-subscriber", @@ -3195,7 +3048,7 @@ dependencies = [ "serde", "serde_json", "sqlx", - "strum 0.26.3", + "strum", "thiserror 2.0.18", "time", "tracing", @@ -3289,7 +3142,7 @@ dependencies = [ "derive_more", "log", "new_debug_unreachable", - "phf 0.13.1", + "phf", "phf_codegen", "precomputed-hash", "rustc-hash", @@ -3763,7 +3616,7 @@ checksum = "a18596f8c785a729f2819c0f6a7eae6ebeebdfffbfe4214ae6b087f690e31901" dependencies = [ "new_debug_unreachable", "parking_lot", - "phf_shared 0.13.1", + "phf_shared", "precomputed-hash", ] @@ -3774,7 +3627,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "585635e46db231059f76c5849798146164652513eb9e8ab2685939dd90f29b69" dependencies = [ "phf_generator", - "phf_shared 0.13.1", + "phf_shared", "proc-macro2", "quote", ] @@ -3790,39 +3643,12 @@ dependencies = [ "unicode-properties", ] -[[package]] -name = "strsim" -version = "0.11.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7da8b5736845d9f2fcb837ea5d9e2628564b3b043a70948a3f0b778838c5fb4f" - [[package]] name = "strum" version = "0.26.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8fec0f0aef304996cf250b31b5a10dee7980c85da9d759361292b8bca5a18f06" -[[package]] -name = "strum" -version = "0.27.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "af23d6f6c1a224baef9d3f61e287d2761385a5b88fdab4eb4c6f11aeb54c4bcf" -dependencies = [ - "strum_macros", -] - -[[package]] -name = "strum_macros" -version = "0.27.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7695ce3845ea4b33927c055a39dc438a45b059f7c1b3d91d38d10355fb8cbca7" -dependencies = [ - "heck 0.5.0", - "proc-macro2", - "quote", - "syn 2.0.117", -] - [[package]] name = "subtle" version = "2.6.1" @@ -4053,22 +3879,6 @@ dependencies = [ "windows-sys 0.61.2", ] -[[package]] -name = "tokio-cron-scheduler" -version = "0.15.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1f50e41f200fd8ed426489bd356910ede4f053e30cebfbd59ef0f856f0d7432a" -dependencies = [ - "chrono", - "chrono-tz", - "croner", - "num-derive", - "num-traits", - "tokio", - "tracing", - "uuid", -] - [[package]] name = "tokio-macros" version = "2.6.1" @@ -4716,7 +4526,7 @@ version = "0.2.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "57a9779e9f04d2ac1ce317aee707aa2f6b773afba7b931222bff6983843b1576" dependencies = [ - "phf 0.13.1", + "phf", "phf_codegen", "string_cache", "string_cache_codegen", diff --git a/Cargo.toml b/Cargo.toml index f07044a..449b0fa 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -26,13 +26,11 @@ sea-orm = { version = "1.1.19", features = ["sqlx-postgres", "runtime-tokio-rust uuid = { version = "1.10.0", features = ["v4", "serde"] } chrono = { version = "0.4", features = ["serde"] } -bytes = "1.11.0" futures = "0.3" reqwest = { version = "0.12.28", features = ["json", "stream", "multipart"] } http = "1.4.0" async-trait = "0.1.89" regex = "1.12.2" -infer = "0.19.0" urlencoding = "2.1" url = "2.5.8" @@ -48,9 +46,6 @@ redis = { version = "0.32.7", features = ["tokio-rustls-comp", "safe_iterators"] thiserror = "2.0.18" config = { version = "0.15.19", features = ["toml"] } -tokio-cron-scheduler = "0.15.1" -hex = "0.4.3" -sha2 = "0.10.9" utoipa = { version = "5.0", features = ["axum_extras"] } utoipa-swagger-ui = { version = "9.0", features = ["axum"] } diff --git a/src/application/anime/use_cases.rs b/src/application/anime/use_cases.rs index 5c9fa99..0e8b634 100644 --- a/src/application/anime/use_cases.rs +++ b/src/application/anime/use_cases.rs @@ -3,18 +3,12 @@ //! Orchestrates repository fetching, caching, and image poster processing. //! Returns pure domain types โ€” no DTOs. -use std::sync::Arc; - use deadpool_redis::Pool; -use sea_orm::DatabaseConnection; use crate::domain::entity::anime::*; use crate::domain::error::*; use crate::infrastructure::cache::redis::Cache; use crate::infrastructure::repository::OtakudesuRepository; -use crate::infrastructure::services::images::cache::{ - cache_image_urls_batch_lazy, get_cached_or_original, -}; const INDEX_CACHE_TTL: u64 = 10; const GENRE_LIST_CACHE_TTL: u64 = 3600; @@ -23,22 +17,13 @@ const DEFAULT_CACHE_TTL: u64 = 300; pub struct AnimeUseCases { repository: OtakudesuRepository, redis_pool: Pool, - db: Arc, - semaphore: Option>, } impl AnimeUseCases { - pub fn new( - repository: OtakudesuRepository, - redis_pool: Pool, - db: Arc, - semaphore: Option>, - ) -> Self { + pub fn new(repository: OtakudesuRepository, redis_pool: Pool) -> Self { Self { repository, redis_pool, - db, - semaphore, } } @@ -49,7 +34,7 @@ impl AnimeUseCases { pub async fn get_anime_index(&self) -> Result { self.cache() .get_or_set("anime:index:v2", INDEX_CACHE_TTL, || async { - let mut data = self + let data = self .repository .fetch_anime_index() .await @@ -59,33 +44,6 @@ impl AnimeUseCases { return Err("Empty anime index โ€” refusing to cache".to_string()); } - let mut posters: Vec = data - .ongoing_anime - .iter() - .map(|item| item.poster.clone()) - .collect(); - posters.extend(data.complete_anime.iter().map(|item| item.poster.clone())); - - let cached_posters = cache_image_urls_batch_lazy( - self.db.clone(), - &self.redis_pool, - posters, - self.semaphore.clone(), - ) - .await; - - let ongoing_len = data.ongoing_anime.len(); - for (i, item) in data.ongoing_anime.iter_mut().enumerate() { - if let Some(url) = cached_posters.get(i) { - item.poster = url.clone(); - } - } - for (i, item) in data.complete_anime.iter_mut().enumerate() { - if let Some(url) = cached_posters.get(ongoing_len + i) { - item.poster = url.clone(); - } - } - Ok(data) }) .await @@ -108,39 +66,12 @@ impl AnimeUseCases { let cache_key = format!("anime:detail:{}", slug); self.cache() .get_or_set(&cache_key, DEFAULT_CACHE_TTL, || async { - let mut data = self + let data = self .repository .fetch_anime_detail(&slug) .await .map_err(|e| e.to_string())?; - data.poster = get_cached_or_original( - self.db.clone(), - &self.redis_pool, - &data.poster, - self.semaphore.clone(), - ) - .await; - - let rec_posters: Vec = data - .recommendations - .iter() - .map(|r| r.poster.clone()) - .collect(); - let cached_rec_posters = cache_image_urls_batch_lazy( - self.db.clone(), - &self.redis_pool, - rec_posters, - self.semaphore.clone(), - ) - .await; - - for (i, rec) in data.recommendations.iter_mut().enumerate() { - if let Some(url) = cached_rec_posters.get(i) { - rec.poster = url.clone(); - } - } - Ok(data) }) .await diff --git a/src/application/anime2/use_cases.rs b/src/application/anime2/use_cases.rs index d54e5fb..7fb4314 100644 --- a/src/application/anime2/use_cases.rs +++ b/src/application/anime2/use_cases.rs @@ -8,23 +8,17 @@ //! TODO: Once parsers return domain types, replace shared types with //! `crate::domain::entity::anime::{GenreAnimeItem, SearchAnimeItem, LatestAnimeItem}`. -use std::sync::Arc; - use deadpool_redis::Pool; -use sea_orm::DatabaseConnection; use crate::domain::error::*; use crate::domain::repository::ScrapingRepository; use crate::infrastructure::cache::redis::Cache; use crate::infrastructure::repository::AlqanimeRepository; -use crate::infrastructure::services::images::cache::{ - apply_cached_posters, cache_image_urls_batch_lazy, get_cached_or_original, -}; use crate::infrastructure::repository::parsers::alqanime_parser as parser; use crate::domain::entity::anime::{ - CompleteAnimeItem, FilterAnimeItem, Genre, GenreAnimeItem, HasPoster, LatestAnimeItem, + CompleteAnimeItem, FilterAnimeItem, Genre, GenreAnimeItem, LatestAnimeItem, OngoingAnimeItemWithScore, Pagination, SearchAnimeItem, }; @@ -56,15 +50,6 @@ pub struct Anime2Item { pub anime_url: String, } -impl HasPoster for Anime2Item { - fn poster(&self) -> &str { - &self.poster - } - fn set_poster(&mut self, url: String) { - self.poster = url; - } -} - #[derive(Serialize, Deserialize, Debug, Clone, ToSchema)] pub struct Anime2Data { pub ongoing_anime: Vec, @@ -123,22 +108,13 @@ const COMPLETE_CACHE_TTL: u64 = 300; pub struct Anime2UseCases { repository: AlqanimeRepository, redis_pool: Pool, - db: Arc, - semaphore: Option>, } impl Anime2UseCases { - pub fn new( - repository: AlqanimeRepository, - redis_pool: Pool, - db: Arc, - semaphore: Option>, - ) -> Self { + pub fn new(repository: AlqanimeRepository, redis_pool: Pool) -> Self { Self { repository, redis_pool, - db, - semaphore, } } @@ -168,7 +144,7 @@ impl Anime2UseCases { }) .await .map_err(|e| e.to_string())??; - let mut ongoing: Vec = data + let ongoing: Vec = data .0 .into_iter() .map(|item| Anime2Item { @@ -182,7 +158,7 @@ impl Anime2UseCases { }) .collect(); - let mut complete: Vec = data + let complete: Vec = data .1 .into_iter() .map(|item| Anime2Item { @@ -196,30 +172,6 @@ impl Anime2UseCases { }) .collect(); - let mut posters: Vec = - ongoing.iter().map(|item| item.poster.clone()).collect(); - posters.extend(complete.iter().map(|item| item.poster.clone())); - - let cached_posters = cache_image_urls_batch_lazy( - self.db.clone(), - &self.redis_pool, - posters, - self.semaphore.clone(), - ) - .await; - - let ongoing_len = ongoing.len(); - for (i, item) in ongoing.iter_mut().enumerate() { - if let Some(url) = cached_posters.get(i) { - item.poster = url.clone(); - } - } - for (i, item) in complete.iter_mut().enumerate() { - if let Some(url) = cached_posters.get(ongoing_len + i) { - item.poster = url.clone(); - } - } - Ok(Anime2Response { status: "Ok".to_string(), data: Anime2Data { @@ -293,20 +245,12 @@ impl Anime2UseCases { .fetch_html(&url) .await .map_err(|e| e.to_string())?; - let (mut data, pagination) = tokio::task::spawn_blocking(move || { + let (data, pagination) = tokio::task::spawn_blocking(move || { parser::parse_filter_page(&html, page).map_err(|e| e.to_string()) }) .await .map_err(|e| e.to_string())??; - apply_cached_posters( - &mut data, - self.db.clone(), - &self.redis_pool, - self.semaphore.clone(), - ) - .await; - Ok(FilterResponse { success: true, data, @@ -374,29 +318,6 @@ impl Anime2UseCases { } } - data.poster = get_cached_or_original( - self.db.clone(), - &self.redis_pool, - &data.poster, - self.semaphore.clone(), - ) - .await; - data.poster2 = get_cached_or_original( - self.db.clone(), - &self.redis_pool, - &data.poster2, - self.semaphore.clone(), - ) - .await; - - apply_cached_posters( - &mut data.recommendations, - self.db.clone(), - &self.redis_pool, - self.semaphore.clone(), - ) - .await; - Ok(DetailResponse { status: "Ok".to_string(), data, @@ -420,20 +341,12 @@ impl Anime2UseCases { .fetch_html(&self.repository.genre_page_url(&genre_slug, page)) .await .map_err(|e| e.to_string())?; - let (mut data, _pagination) = tokio::task::spawn_blocking(move || { + let (data, _pagination) = tokio::task::spawn_blocking(move || { parser::parse_genre_page(&html, page).map_err(|e| e.to_string()) }) .await .map_err(|e| e.to_string())??; - apply_cached_posters( - &mut data, - self.db.clone(), - &self.redis_pool, - self.semaphore.clone(), - ) - .await; - Ok(ApiResponse::success(data)) }) .await @@ -454,20 +367,12 @@ impl Anime2UseCases { .fetch_html(&self.repository.search_url(&query, page)) .await .map_err(|e| e.to_string())?; - let (mut data, _pagination) = tokio::task::spawn_blocking(move || { + let (data, _pagination) = tokio::task::spawn_blocking(move || { parser::parse_search_page(&html, page).map_err(|e| e.to_string()) }) .await .map_err(|e| e.to_string())??; - apply_cached_posters( - &mut data, - self.db.clone(), - &self.redis_pool, - self.semaphore.clone(), - ) - .await; - Ok(ApiResponse::success(data)) }) .await @@ -487,20 +392,12 @@ impl Anime2UseCases { .fetch_html(&self.repository.latest_url(page)) .await .map_err(|e| e.to_string())?; - let (mut data, _pagination) = tokio::task::spawn_blocking(move || { + let (data, _pagination) = tokio::task::spawn_blocking(move || { parser::parse_latest_page(&html, page).map_err(|e| e.to_string()) }) .await .map_err(|e| e.to_string())??; - apply_cached_posters( - &mut data, - self.db.clone(), - &self.redis_pool, - self.semaphore.clone(), - ) - .await; - Ok(ApiResponse::success(data)) }) .await @@ -520,20 +417,12 @@ impl Anime2UseCases { .fetch_html(&self.repository.ongoing_url(page)) .await .map_err(|e| e.to_string())?; - let (mut data, _pagination) = tokio::task::spawn_blocking(move || { + let (data, _pagination) = tokio::task::spawn_blocking(move || { parser::parse_ongoing_page(&html, page).map_err(|e| e.to_string()) }) .await .map_err(|e| e.to_string())??; - apply_cached_posters( - &mut data, - self.db.clone(), - &self.redis_pool, - self.semaphore.clone(), - ) - .await; - Ok(ApiResponse::success(data)) }) .await @@ -553,20 +442,12 @@ impl Anime2UseCases { .fetch_html(&self.repository.complete_url(page)) .await .map_err(|e| e.to_string())?; - let (mut data, _pagination) = tokio::task::spawn_blocking(move || { + let (data, _pagination) = tokio::task::spawn_blocking(move || { parser::parse_complete_page(&html, page).map_err(|e| e.to_string()) }) .await .map_err(|e| e.to_string())??; - apply_cached_posters( - &mut data, - self.db.clone(), - &self.redis_pool, - self.semaphore.clone(), - ) - .await; - Ok(ApiResponse::success(data)) }) .await diff --git a/src/application/komik/use_cases.rs b/src/application/komik/use_cases.rs index b6addb4..88bb5ec 100644 --- a/src/application/komik/use_cases.rs +++ b/src/application/komik/use_cases.rs @@ -6,10 +6,7 @@ //! `crate::infrastructure::repository::parsers::komik_parser`. //! TODO: Move response DTOs to `crate::presentation::dto::komik`. -use std::sync::Arc; - use deadpool_redis::Pool; -use sea_orm::DatabaseConnection; use crate::domain::entity::anime::Pagination; use crate::domain::entity::komik::{ChapterData, DetailData, KomikGenre, KomikItem}; @@ -17,9 +14,6 @@ use crate::domain::error::*; use crate::domain::repository::ScrapingRepository; use crate::infrastructure::cache::redis::Cache; use crate::infrastructure::repository::KomikRepository; -use crate::infrastructure::services::images::cache::{ - apply_cached_posters, cache_image_urls_batch_lazy, get_cached_or_original, -}; use crate::infrastructure::repository::parsers::komik_parser as parser; @@ -36,22 +30,13 @@ const SEARCH_CACHE_TTL: u64 = 300; pub struct KomikUseCases { repository: KomikRepository, redis_pool: Pool, - db: Arc, - semaphore: Option>, } impl KomikUseCases { - pub fn new( - repository: KomikRepository, - redis_pool: Pool, - db: Arc, - semaphore: Option>, - ) -> Self { + pub fn new(repository: KomikRepository, redis_pool: Pool) -> Self { Self { repository, redis_pool, - db, - semaphore, } } @@ -94,20 +79,12 @@ impl KomikUseCases { .await .map_err(|e| e.to_string())?; - let (mut komik_list, pagination) = + let (komik_list, pagination) = tokio::task::spawn_blocking(move || parser::parse_genre_page(&html, page)) .await .map_err(|e| e.to_string())? .map_err(|e| e.to_string())?; - apply_cached_posters( - &mut komik_list, - self.db.clone(), - &self.redis_pool, - self.semaphore.clone(), - ) - .await; - Ok((komik_list, pagination)) }) .await @@ -130,20 +107,12 @@ impl KomikUseCases { .await .map_err(|e| e.to_string())?; - let (mut komik_list, pagination) = + let (komik_list, pagination) = tokio::task::spawn_blocking(move || parser::parse_genre_page(&html, page)) .await .map_err(|e| e.to_string())? .map_err(|e| e.to_string())?; - apply_cached_posters( - &mut komik_list, - self.db.clone(), - &self.redis_pool, - self.semaphore.clone(), - ) - .await; - Ok((komik_list, pagination)) }) .await @@ -162,22 +131,12 @@ impl KomikUseCases { .await .map_err(|e| e.to_string())?; - let mut data = + let data = tokio::task::spawn_blocking(move || parser::parse_komik_detail_document(&html)) .await .map_err(|e| e.to_string())? .map_err(|e| e.to_string())?; - if !data.poster.is_empty() { - data.poster = get_cached_or_original( - self.db.clone(), - &self.redis_pool, - &data.poster, - self.semaphore.clone(), - ) - .await; - } - Ok(data) }) .await @@ -196,7 +155,7 @@ impl KomikUseCases { .await .map_err(|e| e.to_string())?; - let mut data = tokio::task::spawn_blocking({ + let data = tokio::task::spawn_blocking({ let chapter_url = chapter_url.clone(); move || parser::parse_komik_chapter_document(&html, &chapter_url) }) @@ -204,14 +163,6 @@ impl KomikUseCases { .map_err(|e| e.to_string())? .map_err(|e| e.to_string())?; - data.images = cache_image_urls_batch_lazy( - self.db.clone(), - &self.redis_pool, - data.images, - self.semaphore.clone(), - ) - .await; - Ok(data) }) .await @@ -278,7 +229,7 @@ impl KomikUseCases { .await .map_err(|e| e.to_string())?; - let (mut komik_list, pagination) = + let (komik_list, pagination) = tokio::task::spawn_blocking(move || parser::parse_genre_page(&html, page)) .await .map_err(|e| e.to_string())? @@ -288,14 +239,6 @@ impl KomikUseCases { return Err(format!("Empty komik {} page {}", list_name, page)); } - apply_cached_posters( - &mut komik_list, - self.db.clone(), - &self.redis_pool, - self.semaphore.clone(), - ) - .await; - Ok((komik_list, pagination)) }) .await @@ -318,20 +261,12 @@ impl KomikUseCases { .await .map_err(|e| e.to_string())?; - let (mut komik_list, pagination) = + let (komik_list, pagination) = tokio::task::spawn_blocking(move || parser::parse_genre_page(&html, page)) .await .map_err(|e| e.to_string())? .map_err(|e| e.to_string())?; - apply_cached_posters( - &mut komik_list, - self.db.clone(), - &self.redis_pool, - self.semaphore.clone(), - ) - .await; - Ok((komik_list, pagination)) }) .await @@ -354,20 +289,12 @@ impl KomikUseCases { .await .map_err(|e| e.to_string())?; - let (mut komik_list, pagination) = + let (komik_list, pagination) = tokio::task::spawn_blocking(move || parser::parse_genre_page(&html, page)) .await .map_err(|e| e.to_string())? .map_err(|e| e.to_string())?; - apply_cached_posters( - &mut komik_list, - self.db.clone(), - &self.redis_pool, - self.semaphore.clone(), - ) - .await; - Ok((komik_list, pagination)) }) .await diff --git a/src/application/proxy/use_cases.rs b/src/application/proxy/use_cases.rs index 8bbca56..1f09306 100644 --- a/src/application/proxy/use_cases.rs +++ b/src/application/proxy/use_cases.rs @@ -1,236 +1 @@ //! Proxy application use cases. -//! -//! Provides proxy fetch, image caching, and audit/repair operations. -//! -//! TODO: Move result types to `crate::presentation::dto::proxy`. -//! TODO: Add event bus integration for ImageRepaired events. -//! TODO: Replace `reqwest::Client::new()` with shared HTTP client from infrastructure. - -use std::sync::Arc; - -use axum::http::StatusCode; -use axum::response::Response; -use serde::Serialize; -use tracing::{error, info, warn}; -use utoipa::ToSchema; - -use crate::domain::error::*; -use crate::domain::repository::ImageCacheRepository; -use crate::infrastructure::repository::{ProxyRepository, SeaOrmImageCacheRepository}; -use crate::infrastructure::services::images::cache::ImageCache; - -// ============================================================================ -// Result types โ€” TEMPORARY: move to presentation layer -// ============================================================================ - -/// Result of caching an image URL. -/// TODO: Move to presentation::dto::proxy -#[derive(Debug, Serialize, ToSchema)] -pub struct ImageCacheResult { - pub success: bool, - pub original_url: String, - pub cdn_url: String, - pub from_cache: bool, - pub pending: Option, -} - -/// Result of auditing/repairing a cached image URL. -/// TODO: Move to presentation::dto::proxy -#[derive(Debug, Serialize, ToSchema)] -pub struct AuditImageCacheResult { - pub success: bool, - pub original_url: String, - pub cdn_url: Option, - pub was_accessible: bool, - pub re_uploaded: bool, - pub message: String, -} - -// ============================================================================ -// Use case struct -// ============================================================================ - -pub struct ProxyUseCases { - repository: ProxyRepository, - image_cache_repo: Arc, -} - -impl ProxyUseCases { - pub fn new( - repository: ProxyRepository, - image_cache_repo: Arc, - ) -> Self { - Self { - repository, - image_cache_repo, - } - } - - fn build_image_cache(&self) -> ImageCache { - ImageCache::new(self.image_cache_repo.clone()) - } - - pub async fn fetch_with_proxy_only(&self, url: String) -> Result { - let fetch_result = self - .repository - .fetch_with_proxy_url(&url) - .await - .map_err(|e| DomainError::Scraping(ScrapingError::Http(e.to_string())))?; - - let mut builder = Response::builder().status(StatusCode::OK); - if let Some(content_type) = fetch_result.content_type { - builder = builder.header("Content-Type", content_type); - } - - builder - .body(fetch_result.data.into()) - .map_err(|e| DomainError::Repository(RepositoryError::Network(e.to_string()))) - } - - pub async fn image_cache( - &self, - url: String, - lazy: bool, - ) -> Result { - let cache = self.build_image_cache(); - - if let Some(cdn_url) = cache.get_cdn_url(&url).await { - return Ok(ImageCacheResult { - success: true, - original_url: url, - cdn_url, - from_cache: true, - pending: None, - }); - } - - if lazy { - let repo = self.image_cache_repo.clone(); - let url_clone = url.clone(); - tokio::spawn(async move { - let cache = ImageCache::new(repo); - match cache.get_or_cache(&url_clone).await { - Ok(cdn) => info!("[LazyCache] Cached {} -> {}", url_clone, cdn), - Err(e) => warn!("[LazyCache] Failed {}: {}", url_clone, e), - } - }); - return Ok(ImageCacheResult { - success: true, - original_url: url.clone(), - cdn_url: url, - from_cache: false, - pending: Some(true), - }); - } - - match cache.get_or_cache(&url).await { - Ok(cdn_url) => Ok(ImageCacheResult { - success: true, - original_url: url, - cdn_url, - from_cache: false, - pending: None, - }), - Err(e) => { - error!("ImageCache error: {}", e); - Ok(ImageCacheResult { - success: false, - original_url: url.clone(), - cdn_url: url, - from_cache: false, - pending: None, - }) - } - } - } - - pub async fn audit_image_cache( - &self, - url: String, - ) -> Result { - let cache = self.build_image_cache(); - let mut cdn_opt = cache.get_cdn_url(&url).await; - let mut original = url.clone(); - - if cdn_opt.is_none() { - if let Some(orig) = cache.find_original_from_cdn(&url).await { - info!("SmartAudit: {} recognized as CDN, original {}", url, orig); - original = orig; - cdn_opt = Some(url.clone()); - } - } - - if let Some(cdn_url) = cdn_opt { - let client = reqwest::Client::new(); - let mut accessible = false; - - match client.get(&cdn_url).send().await { - Ok(resp) if resp.status().is_success() => { - if let Ok(bytes) = resp.bytes().await { - if infer::get(&bytes) - .map(|k| k.mime_type().starts_with("image/")) - .unwrap_or(false) - { - accessible = true; - } else { - warn!("CDN {} returned non-image content", cdn_url); - } - } - } - Ok(resp) => warn!("CDN {} status {}", cdn_url, resp.status()), - Err(e) => warn!("CDN {} fetch error {}", cdn_url, e), - } - - if accessible { - return Ok(AuditImageCacheResult { - success: true, - original_url: original, - cdn_url: Some(cdn_url), - was_accessible: true, - re_uploaded: false, - message: "CDN URL is accessible and the image is valid".to_string(), - }); - } - - info!("CDN {} inaccessible, purging and reuploading", cdn_url); - let _ = cache.invalidate(&original).await; - match cache.get_or_cache(&original).await { - Ok(new_cdn) => Ok(AuditImageCacheResult { - success: true, - original_url: original, - cdn_url: Some(new_cdn), - was_accessible: false, - re_uploaded: true, - message: "CDN URL was inaccessible, re-uploaded".to_string(), - }), - Err(e) => Ok(AuditImageCacheResult { - success: false, - original_url: original, - cdn_url: None, - was_accessible: false, - re_uploaded: false, - message: format!("Re-upload failed: {}", e), - }), - } - } else { - match cache.get_or_cache(&original).await { - Ok(new_cdn) => Ok(AuditImageCacheResult { - success: true, - original_url: original, - cdn_url: Some(new_cdn), - was_accessible: false, - re_uploaded: true, - message: "Cached newly".to_string(), - }), - Err(e) => Ok(AuditImageCacheResult { - success: false, - original_url: original, - cdn_url: None, - was_accessible: false, - re_uploaded: false, - message: format!("Cache failed: {}", e), - }), - } - } - } -} diff --git a/src/bootstrap/mod.rs b/src/bootstrap/mod.rs index 0a986bd..fd0a94e 100644 --- a/src/bootstrap/mod.rs +++ b/src/bootstrap/mod.rs @@ -77,26 +77,15 @@ impl Application { // App State components let db_arc = Arc::new(db); - let image_processing_semaphore = Arc::new(tokio::sync::Semaphore::new( - CONFIG.image_processing_concurrency, - )); let event_bus = Arc::new(crate::events::bus::EventBus::new()); let redis_pool = crate::infrastructure::cache::redis_pool::redis_pool() .map_err(|e| anyhow::anyhow!("Failed to init Redis pool: {}", e))?; - use crate::infrastructure::repository::SeaOrmImageCacheRepository; - let image_cache_repo = Arc::new(SeaOrmImageCacheRepository::new( - db_arc.clone(), - redis_pool.clone(), - )); - let app_state = Arc::new(AppState { redis_pool, db: db_arc.clone(), - image_processing_semaphore, event_bus: event_bus.clone(), - image_cache_repo, }); let app = crate::presentation::router::build_router(app_state.clone())?; diff --git a/src/bootstrap/setup.rs b/src/bootstrap/setup.rs index aeecd2d..a40d98b 100644 --- a/src/bootstrap/setup.rs +++ b/src/bootstrap/setup.rs @@ -1,47 +1,9 @@ //! Database schema initialization. -use sea_orm::{ConnectionTrait, DatabaseConnection, Schema, Statement}; +use sea_orm::DatabaseConnection; use tracing::info; -use crate::infrastructure::persistence::entities::image_cache; - -pub async fn init(db: &DatabaseConnection) -> Result<(), sea_orm::DbErr> { - info!("๐Ÿš€ Initializing database schema..."); - let backend = db.get_database_backend(); - let schema = Schema::new(backend); - - let tables = vec![( - "ImageCache", - schema - .create_table_from_entity(image_cache::Entity) - .if_not_exists() - .to_owned(), - )]; - - for (name, stmt) in tables { - match db.execute(backend.build(&stmt)).await { - Ok(_) => info!(" โœ“ Table '{}' checked/created", name), - Err(e) => { - tracing::error!(" [!] Failed to create table '{}': {}", name, e); - return Err(e); - } - } - } - - let index_sql = - "CREATE INDEX IF NOT EXISTS idx_image_cache_cdn_url ON \"ImageCache\" (cdn_url)"; - match db.execute(Statement::from_string(backend, index_sql)).await { - Ok(_) => info!(" โœ“ Index 'idx_image_cache_cdn_url' ensured"), - Err(e) => { - let err_str = e.to_string(); - if err_str.contains("already exists") || err_str.contains("duplicate") { - info!(" โœ“ Index 'idx_image_cache_cdn_url' already exists"); - } else { - tracing::error!(" [!] Failed to create index on ImageCache: {}", e); - } - } - } - - info!("โœ… Database schema initialization complete."); +pub async fn init(_db: &DatabaseConnection) -> Result<(), sea_orm::DbErr> { + info!("โœ… Database schema initialization complete (no tables to create)."); Ok(()) } diff --git a/src/config/mod.rs b/src/config/mod.rs index 79eb3c1..7e7b62b 100644 --- a/src/config/mod.rs +++ b/src/config/mod.rs @@ -47,10 +47,6 @@ pub struct AppConfig { #[serde(default)] pub db: DbConfig, - /// Max concurrent image processing tasks - #[serde(default = "default_image_processing_concurrency")] - pub image_processing_concurrency: usize, - /// Domain/URL values that may change between deployments #[serde(default)] pub urls: UrlConfig, @@ -60,18 +56,12 @@ pub struct AppConfig { pub struct UrlConfig { #[serde(default = "default_site_url")] pub site_url: String, - #[serde(default = "default_picser_api_url")] - pub picser_api_url: String, - #[serde(default = "default_fallback_upload_api_url")] - pub fallback_upload_api_url: String, } impl Default for UrlConfig { fn default() -> Self { Self { site_url: default_site_url(), - picser_api_url: default_picser_api_url(), - fallback_upload_api_url: default_fallback_upload_api_url(), } } } @@ -181,22 +171,10 @@ fn default_log_level() -> String { "info".to_string() } -fn default_image_processing_concurrency() -> usize { - 5 -} - fn default_site_url() -> String { "https://asepharyana.my.id".to_string() } -fn default_picser_api_url() -> String { - "https://picser.asepharyana.my.id/api/upload".to_string() -} - -fn default_fallback_upload_api_url() -> String { - "https://upload.asepharyana.my.id/api/upload".to_string() -} - fn default_db_max_connections() -> u32 { 100 } diff --git a/src/domain/entity/anime.rs b/src/domain/entity/anime.rs index 9924691..54bf50d 100644 --- a/src/domain/entity/anime.rs +++ b/src/domain/entity/anime.rs @@ -290,115 +290,4 @@ pub struct FilterAnimeItem { pub anime_url: String, } -// ============================================================================ -// TRAITS -// ============================================================================ -/// Trait for types that have a poster image URL -pub trait HasPoster { - fn poster(&self) -> &str; - fn set_poster(&mut self, url: String); -} - -// ============================================================================ -// TRAIT IMPLEMENTATIONS -// ============================================================================ - -impl HasPoster for OngoingAnimeItem { - fn poster(&self) -> &str { - &self.poster - } - fn set_poster(&mut self, url: String) { - self.poster = url; - } -} - -impl HasPoster for CompleteAnimeItem { - fn poster(&self) -> &str { - &self.poster - } - fn set_poster(&mut self, url: String) { - self.poster = url; - } -} - -impl HasPoster for CompleteAnimeListItem { - fn poster(&self) -> &str { - &self.poster - } - fn set_poster(&mut self, url: String) { - self.poster = url; - } -} - -impl HasPoster for OngoingAnimeListItem { - fn poster(&self) -> &str { - &self.poster - } - fn set_poster(&mut self, url: String) { - self.poster = url; - } -} - -impl HasPoster for LatestAnimeItem { - fn poster(&self) -> &str { - &self.poster - } - fn set_poster(&mut self, url: String) { - self.poster = url; - } -} - -impl HasPoster for SearchAnimeItem { - fn poster(&self) -> &str { - &self.poster - } - fn set_poster(&mut self, url: String) { - self.poster = url; - } -} - -impl HasPoster for GenreAnimeItem { - fn poster(&self) -> &str { - &self.poster - } - fn set_poster(&mut self, url: String) { - self.poster = url; - } -} - -impl HasPoster for Recommendation { - fn poster(&self) -> &str { - &self.poster - } - fn set_poster(&mut self, url: String) { - self.poster = url; - } -} - -impl HasPoster for Anime2Item { - fn poster(&self) -> &str { - &self.poster - } - fn set_poster(&mut self, url: String) { - self.poster = url; - } -} - -impl HasPoster for OngoingAnimeItemWithScore { - fn poster(&self) -> &str { - &self.poster - } - fn set_poster(&mut self, url: String) { - self.poster = url; - } -} - -impl HasPoster for FilterAnimeItem { - fn poster(&self) -> &str { - &self.poster - } - fn set_poster(&mut self, url: String) { - self.poster = url; - } -} diff --git a/src/domain/entity/komik.rs b/src/domain/entity/komik.rs index dc876c8..614a3e2 100644 --- a/src/domain/entity/komik.rs +++ b/src/domain/entity/komik.rs @@ -3,8 +3,6 @@ use serde::{Deserialize, Serialize}; use utoipa::ToSchema; -use crate::domain::entity::anime::HasPoster; - #[derive(Serialize, Deserialize, Debug, Clone, ToSchema)] pub struct KomikGenre { pub name: String, @@ -53,12 +51,3 @@ pub struct KomikItem { pub r#type: String, pub komik_url: String, } - -impl HasPoster for KomikItem { - fn poster(&self) -> &str { - &self.poster - } - fn set_poster(&mut self, url: String) { - self.poster = url; - } -} diff --git a/src/domain/repository/image_cache.rs b/src/domain/repository/image_cache.rs deleted file mode 100644 index 3f1e2db..0000000 --- a/src/domain/repository/image_cache.rs +++ /dev/null @@ -1,19 +0,0 @@ -//! Repository trait for image cache (original โ†’ CDN mapping). - -use async_trait::async_trait; - -/// Repository for managing cached image URL mappings. -#[async_trait] -pub trait ImageCacheRepository: Send + Sync { - async fn get_from_redis(&self, key: &str) -> Option; - async fn set_in_redis(&self, key: &str, value: &str, ttl: u64) -> Result<(), String>; - async fn get_from_db(&self, original_url: &str) -> Result, String>; - async fn save_to_db(&self, original_url: &str, cdn_url: &str) -> Result<(), String>; - async fn find_original_from_cdn(&self, cdn_url: &str) -> Result, String>; - async fn delete_from_db(&self, original_url: &str) -> Result<(), String>; - async fn delete_from_redis(&self, key: &str) -> Result<(), String>; - async fn get_lock(&self, key: &str) -> bool; - async fn set_lock(&self, key: &str, ttl: u64) -> Result<(), String>; - async fn release_lock(&self, key: &str) -> Result<(), String>; - async fn invalidate_api_caches(&self, patterns: Vec<&str>) -> Result<(), String>; -} diff --git a/src/domain/repository/mod.rs b/src/domain/repository/mod.rs index 28d19c1..76b16b4 100644 --- a/src/domain/repository/mod.rs +++ b/src/domain/repository/mod.rs @@ -1,5 +1,3 @@ -pub mod image_cache; pub mod scraping; -pub use image_cache::ImageCacheRepository; pub use scraping::ScrapingRepository; diff --git a/src/events/bus.rs b/src/events/bus.rs index 2a95daa..f09905c 100644 --- a/src/events/bus.rs +++ b/src/events/bus.rs @@ -143,13 +143,4 @@ impl Event for OrderCreated { const NAME: &'static str = "order.created"; } -/// Image repaired event. -#[derive(Clone, Debug, serde::Serialize, serde::Deserialize)] -pub struct ImageRepaired { - pub original_url: String, - pub cdn_url: String, -} -impl Event for ImageRepaired { - const NAME: &'static str = "image.repaired"; -} diff --git a/src/infrastructure/mod.rs b/src/infrastructure/mod.rs index fb4ed2b..07e6d9f 100644 --- a/src/infrastructure/mod.rs +++ b/src/infrastructure/mod.rs @@ -1,6 +1,4 @@ pub mod cache; -pub mod persistence; pub mod repository; pub mod scraping; -pub mod services; pub mod utils; diff --git a/src/infrastructure/persistence/entities/image_cache.rs b/src/infrastructure/persistence/entities/image_cache.rs deleted file mode 100644 index 073fb6a..0000000 --- a/src/infrastructure/persistence/entities/image_cache.rs +++ /dev/null @@ -1,71 +0,0 @@ -//! `SeaORM` Entity for ImageCache - stores URL mappings for CDN caching - -use sea_orm::entity::prelude::*; -use serde::{Deserialize, Serialize}; - -#[derive(Copy, Clone, Default, Debug, DeriveEntity)] -pub struct Entity; - -impl EntityName for Entity { - fn table_name(&self) -> &str { - "ImageCache" - } -} - -#[derive(Clone, Debug, PartialEq, DeriveModel, DeriveActiveModel, Eq, Serialize, Deserialize)] -pub struct Model { - pub id: String, - pub original_url: String, - pub cdn_url: String, - pub created_at: DateTimeUtc, - pub expires_at: Option, -} - -#[derive(Copy, Clone, Debug, EnumIter, DeriveColumn)] -pub enum Column { - Id, - #[sea_orm(column_name = "original_url")] - OriginalUrl, - #[sea_orm(column_name = "cdn_url")] - CdnUrl, - #[sea_orm(column_name = "created_at")] - CreatedAt, - #[sea_orm(column_name = "expires_at")] - ExpiresAt, -} - -#[derive(Copy, Clone, Debug, EnumIter, DerivePrimaryKey)] -pub enum PrimaryKey { - Id, -} - -impl PrimaryKeyTrait for PrimaryKey { - type ValueType = String; - fn auto_increment() -> bool { - false - } -} - -#[derive(Copy, Clone, Debug, EnumIter)] -pub enum Relation {} - -impl ColumnTrait for Column { - type EntityName = Entity; - fn def(&self) -> ColumnDef { - match self { - Self::Id => ColumnType::String(StringLen::N(36u32)).def(), - Self::OriginalUrl => ColumnType::String(StringLen::N(512u32)).def().unique(), - Self::CdnUrl => ColumnType::String(StringLen::N(512u32)).def(), - Self::CreatedAt => ColumnType::TimestampWithTimeZone.def(), - Self::ExpiresAt => ColumnType::TimestampWithTimeZone.def().null(), - } - } -} - -impl RelationTrait for Relation { - fn def(&self) -> RelationDef { - match *self {} - } -} - -impl ActiveModelBehavior for ActiveModel {} diff --git a/src/infrastructure/persistence/entities/mod.rs b/src/infrastructure/persistence/entities/mod.rs deleted file mode 100644 index 44e6ce2..0000000 --- a/src/infrastructure/persistence/entities/mod.rs +++ /dev/null @@ -1 +0,0 @@ -pub mod image_cache; diff --git a/src/infrastructure/persistence/mod.rs b/src/infrastructure/persistence/mod.rs deleted file mode 100644 index 0b8f0b5..0000000 --- a/src/infrastructure/persistence/mod.rs +++ /dev/null @@ -1 +0,0 @@ -pub mod entities; diff --git a/src/infrastructure/repository/image_cache_seaorm.rs b/src/infrastructure/repository/image_cache_seaorm.rs deleted file mode 100644 index 089ec77..0000000 --- a/src/infrastructure/repository/image_cache_seaorm.rs +++ /dev/null @@ -1,131 +0,0 @@ -//! SeaORM-backed implementation of ImageCacheRepository. - -use async_trait::async_trait; -use chrono::Utc; -use deadpool_redis::Pool as RedisPool; -use sea_orm::{ActiveModelTrait, ColumnTrait, DatabaseConnection, EntityTrait, QueryFilter, Set}; -use std::sync::Arc; - -use crate::domain::repository::ImageCacheRepository; -use crate::infrastructure::cache::redis::Cache; -use crate::infrastructure::persistence::entities::image_cache; - -pub struct SeaOrmImageCacheRepository { - db: Arc, - redis: RedisPool, -} - -impl SeaOrmImageCacheRepository { - pub fn new(db: Arc, redis: RedisPool) -> Self { - Self { db, redis } - } -} - -#[async_trait] -impl ImageCacheRepository for SeaOrmImageCacheRepository { - async fn get_from_redis(&self, key: &str) -> Option { - Cache::new(&self.redis).get::(key).await - } - - async fn set_in_redis(&self, key: &str, value: &str, ttl: u64) -> Result<(), String> { - Cache::new(&self.redis) - .set_with_ttl(key, &value, ttl) - .await - .map_err(|e| e.to_string()) - } - - async fn get_from_db(&self, original_url: &str) -> Result, String> { - let entry = image_cache::Entity::find() - .filter(image_cache::Column::OriginalUrl.eq(original_url)) - .one(self.db.as_ref()) - .await - .map_err(|e| e.to_string())?; - Ok(entry.map(|m| m.cdn_url)) - } - - async fn save_to_db(&self, original_url: &str, cdn_url: &str) -> Result<(), String> { - let model = image_cache::ActiveModel { - id: Set(uuid::Uuid::new_v4().to_string()), - original_url: Set(original_url.to_string()), - cdn_url: Set(cdn_url.to_string()), - created_at: Set(Utc::now()), - expires_at: Set(None), - }; - model - .insert(self.db.as_ref()) - .await - .map_err(|e| e.to_string())?; - Ok(()) - } - - async fn find_original_from_cdn(&self, cdn_url: &str) -> Result, String> { - let entry = image_cache::Entity::find() - .filter(image_cache::Column::CdnUrl.eq(cdn_url)) - .one(self.db.as_ref()) - .await - .map_err(|e| e.to_string())?; - Ok(entry.map(|m| m.original_url)) - } - - async fn delete_from_db(&self, original_url: &str) -> Result<(), String> { - image_cache::Entity::delete_many() - .filter(image_cache::Column::OriginalUrl.eq(original_url)) - .exec(self.db.as_ref()) - .await - .map_err(|e| e.to_string())?; - Ok(()) - } - - async fn delete_from_redis(&self, key: &str) -> Result<(), String> { - Cache::new(&self.redis) - .delete(key) - .await - .map_err(|e| e.to_string()) - } - - async fn get_lock(&self, key: &str) -> bool { - Cache::new(&self.redis).get::(key).await.is_some() - } - - async fn set_lock(&self, key: &str, ttl: u64) -> Result<(), String> { - Cache::new(&self.redis) - .set_with_ttl(key, &true, ttl) - .await - .map_err(|e| e.to_string()) - } - - async fn release_lock(&self, key: &str) -> Result<(), String> { - self.delete_from_redis(key).await - } - - async fn invalidate_api_caches(&self, patterns: Vec<&str>) -> Result<(), String> { - use deadpool_redis::redis::{cmd, AsyncCommands}; - - let mut conn = self.redis.get().await.map_err(|e| e.to_string())?; - - for pattern in patterns { - let mut cursor: u64 = 0; - loop { - let (new_cursor, keys): (u64, Vec) = cmd("SCAN") - .arg(cursor) - .arg("MATCH") - .arg(pattern) - .arg("COUNT") - .arg(100) - .query_async(&mut *conn) - .await - .map_err(|e| e.to_string())?; - - if !keys.is_empty() { - let _: usize = conn.del(&keys).await.map_err(|e| e.to_string())?; - } - - cursor = new_cursor; - if cursor == 0 { - break; - } - } - } - Ok(()) - } -} diff --git a/src/infrastructure/repository/mod.rs b/src/infrastructure/repository/mod.rs index e468798..f4e0258 100644 --- a/src/infrastructure/repository/mod.rs +++ b/src/infrastructure/repository/mod.rs @@ -1,12 +1,10 @@ pub mod alqanime; -pub mod image_cache_seaorm; pub mod komik; pub mod otakudesu; pub mod parsers; pub mod proxy; pub use alqanime::AlqanimeRepository; -pub use image_cache_seaorm::SeaOrmImageCacheRepository; pub use komik::KomikRepository; pub use otakudesu::OtakudesuRepository; pub use proxy::ProxyRepository; diff --git a/src/infrastructure/repository/parsers/alqanime_parser.rs b/src/infrastructure/repository/parsers/alqanime_parser.rs index 27199f9..e4f7d78 100644 --- a/src/infrastructure/repository/parsers/alqanime_parser.rs +++ b/src/infrastructure/repository/parsers/alqanime_parser.rs @@ -1,7 +1,7 @@ use crate::domain::entity::anime::{ - CompleteAnimeItem, DetailGenre, FilterAnimeItem, Genre, GenreAnimeItem, HasPoster, - LatestAnimeItem, OngoingAnimeItem, OngoingAnimeItemWithScore, Pagination, - PaginationWithStringPages, SearchAnimeItem, + CompleteAnimeItem, DetailGenre, FilterAnimeItem, Genre, GenreAnimeItem, LatestAnimeItem, + OngoingAnimeItem, OngoingAnimeItemWithScore, Pagination, PaginationWithStringPages, + SearchAnimeItem, }; use crate::domain::error::ScrapingError; use crate::infrastructure::scraping::parsing_utils::parse_html; @@ -41,15 +41,6 @@ pub struct AlqRecommendation { pub r#type: String, } -impl HasPoster for AlqRecommendation { - fn poster(&self) -> &str { - &self.poster - } - fn set_poster(&mut self, url: String) { - self.poster = url; - } -} - use regex::Regex; use scraper::Selector; use std::sync::LazyLock; diff --git a/src/infrastructure/services/images/cache.rs b/src/infrastructure/services/images/cache.rs deleted file mode 100644 index f2e9940..0000000 --- a/src/infrastructure/services/images/cache.rs +++ /dev/null @@ -1,1142 +0,0 @@ -//! Image caching helper using Picser CDN (picser.pages.dev). -//! -//! This module provides utilities to cache images via jsDelivr CDN -//! with database storage for URL mapping. - -use crate::config::CONFIG; -use crate::domain::repository::ImageCacheRepository; -use crate::infrastructure::repository::SeaOrmImageCacheRepository; -use deadpool_redis::Pool as RedisPool; -use reqwest::Client; -use sea_orm::DatabaseConnection; -use serde::{Deserialize, Serialize}; -use std::sync::Arc; -use tracing::{debug, error, warn}; - -use crate::infrastructure::cache::redis::Cache; - -/// Default TTL for image cache in Redis (24 hours) -pub const CACHE_TTL_IMAGE: u64 = 86400; -use crate::infrastructure::utils::http_client::http_client; - -/// Default TTL for image cache in Redis (24 hours) -pub const IMAGE_CACHE_TTL: u64 = CACHE_TTL_IMAGE; - -/// Redis key prefix for image cache -pub const IMAGE_CACHE_PREFIX: &str = "img_cache"; - -/// Redis key prefix for caching locks (to prevent duplicate uploads) -pub const IMAGE_CACHE_LOCK_PREFIX: &str = "img_cache_lock"; - -/// Lock TTL (60 seconds - enough time for upload to complete) -pub const IMAGE_CACHE_LOCK_TTL: u64 = 60; - -/// Static Picser API endpoints in priority order. Configured endpoint is inserted after primary. -pub const STATIC_PICSER_API_ENDPOINTS: &[&str] = &[ - "https://picser-two.vercel.app/api/upload", - "https://picser-mytheclipse8647-ahoqi9ef.leapcell.dev/api/upload", - "https://picser.pages.dev/api/upload", -]; - -/// Create a hash of the URL for cache key -pub fn url_hash(url: &str) -> String { - use sha2::{Digest, Sha256}; - let mut hasher = Sha256::new(); - hasher.update(url.as_bytes()); - let result = hasher.finalize(); - hex::encode(&result[..16]) // Use 16 bytes for collision-resistant key -} - -/// Helper to convert any image URL to a fast WP.com (Jetpack) CDN URL. -/// This acts as a high-speed proxy even before Picser finishes caching. -pub fn to_wp_cdn(url: &str) -> String { - if url.is_empty() { - return url.to_string(); - } - - // If already a CDN URL, return as is - if url.contains("picser.pages.dev") - || url.contains("jsdelivr.net") - || url.contains("wp.com") - || url.contains("imagecdn.app") - { - return url.to_string(); - } - - // Remove protocol for wp.com format - let clean_url = url.trim().replace("https://", "").replace("http://", ""); - - // Use i0, i1, i2 or i3 based on hash to distribute load - let hash = url.len() % 4; - format!("https://i{}.wp.com/{}", hash, clean_url) -} - -/// Response from Picser API (/api/upload) -#[derive(Debug, Clone, Serialize, Deserialize)] -pub struct PicserResponse { - #[serde(default)] - pub success: bool, - pub url: Option, - pub urls: Option, - pub filename: Option, - pub size: Option, - #[serde(rename = "type")] - pub content_type: Option, - pub commit_sha: Option, - pub github_url: Option, - pub error: Option, -} - -#[derive(Debug, Clone, Serialize, Deserialize)] -pub struct PicserUrls { - pub github: Option, - pub raw: Option, - pub jsdelivr: Option, - pub jsdelivr_commit: Option, -} - -/// Response from fallback upload API (/api/upload) -#[derive(Debug, Clone, Serialize, Deserialize)] -pub struct FallbackUploadResponse { - pub download_url: String, -} - -/// Configuration for image cache -#[derive(Debug, Clone)] -pub struct ImageCacheConfig { - /// GitHub token for Picser API (optional - uses public upload if not set) - pub github_token: Option, - /// GitHub owner for uploads - pub github_owner: String, - /// GitHub repo for uploads - pub github_repo: String, - /// GitHub branch - pub github_branch: String, - /// Upload folder - pub folder: String, -} - -impl Default for ImageCacheConfig { - fn default() -> Self { - Self { - github_token: None, - github_owner: "sh20raj".to_string(), - github_repo: "picser".to_string(), - github_branch: "main".to_string(), - folder: "uploads".to_string(), - } - } -} - -/// Image cache service -pub struct ImageCache { - repo: Arc, - client: Client, - _config: ImageCacheConfig, - semaphore: Option>, -} - -// Add imports for Request Coalescing -use dashmap::DashMap; -use std::sync::LazyLock; -use tokio::sync::broadcast; - -// Global In-Flight Uploads Map -// Maps Original URL -> Broadcast Sender -static IN_FLIGHT_UPLOADS: LazyLock>>> = - LazyLock::new(DashMap::new); - -impl ImageCache { - /// Create a new image cache instance - pub fn new(repo: Arc) -> Self { - Self { - repo, - client: http_client().client().clone(), // Reuse global HTTP client for connection pooling - _config: ImageCacheConfig::default(), - semaphore: None, - } - } - - pub fn with_config(repo: Arc, config: ImageCacheConfig) -> Self { - Self { - repo, - client: http_client().client().clone(), // Reuse global HTTP client - _config: config, - semaphore: None, - } - } - - /// Set concurrency limiter - pub fn with_semaphore(mut self, semaphore: std::sync::Arc) -> Self { - self.semaphore = Some(semaphore); - self - } - - /// Get CDN URL for an image, caching if needed - pub async fn get_or_cache(&self, original_url: &str) -> Result { - let cache_key = format!("{}:{}", IMAGE_CACHE_PREFIX, url_hash(original_url)); - let lock_key = format!("{}:{}", IMAGE_CACHE_LOCK_PREFIX, url_hash(original_url)); - - // 1. Check Redis cache first - - if let Some(cached_url) = self.repo.get_from_redis(&cache_key).await { - debug!("ImageCache: Redis hit for {}", original_url); - return Ok(cached_url); - } - - // 2. Check Request Coalescing (SingleFlight) - // This handles concurrent requests in this process/instance - let (tx, is_leader) = { - use dashmap::mapref::entry::Entry; - match IN_FLIGHT_UPLOADS.entry(original_url.to_string()) { - Entry::Occupied(entry) => { - debug!("ImageCache: Joining in-flight upload for {}", original_url); - (entry.get().clone(), false) - } - Entry::Vacant(entry) => { - let (tx, _) = broadcast::channel(1); - entry.insert(tx.clone()); - debug!("ImageCache: Starting leader upload for {}", original_url); - (tx, true) - } - } - }; - - if !is_leader { - // Follower: Wait for result - let mut rx = tx.subscribe(); - return match rx.recv().await { - Ok(Ok(url)) => Ok(url), - Ok(Err(e)) => Err(e), - Err(e) => { - warn!( - "ImageCache: Coalesce receive error for {}: {:?}", - original_url, e - ); - Err("Upload coalescing failed".to_string()) - } - }; - } - - // Leader: Perform the work - // We wrap the work in a closure/block to easily capture the result - let result = async { - // 3. Check database (Double check inside leader to be sure) - if let Some(cached_url) = self.repo.get_from_db(original_url).await? { - // Store in Redis for faster access - let _ = self - .repo - .set_in_redis(&cache_key, &cached_url, IMAGE_CACHE_TTL) - .await; - debug!("ImageCache: DB hit for {}", original_url); - return Ok(cached_url); - } - - // 4. Check if another process is already caching this URL (Distributed Lock check) - if self.repo.get_lock(&lock_key).await { - // Even if locked by another process, strict single-flight within this instance - // is good. But if another process is working, we might want to wait or just return error? - // Current logic returns error. - debug!( - "ImageCache: Already being cached by another process: {}", - original_url - ); - return Err(format!("URL {} is already being cached", original_url)); - } - - // 5. Acquire lock in Redis - let _ = self.repo.set_lock(&lock_key, IMAGE_CACHE_LOCK_TTL).await; - - // 6. Upload - debug!("ImageCache: Miss - uploading {} to Picser", original_url); - - // Acquire permit if semaphore is set - let _permit = if let Some(sem) = &self.semaphore { - match sem.acquire().await { - Ok(p) => Some(p), - Err(e) => { - let _ = self.repo.release_lock(&lock_key).await; - return Err(e.to_string()); - } - } - } else { - None - }; - - // Work - let work_result = async { - // Upload to Picser - let cdn_url = self.upload_to_picser(original_url).await?; - - // 6.5. Verify CDN URL Propagation (Self-Test before caching) - // CDNs like jsDelivr can take a few seconds to propagate after a GitHub commit. - // We retry 3 times with backoff to ensure we only return and cache a functional link. - let mut is_valid = false; - let mut last_verify_error = String::from("Verification not started"); - - for attempt in 1..=10 { - debug!( - "ImageCache: Verifying CDN URL {} (Attempt {})", - cdn_url, attempt - ); - match self.verify_cdn_url(&cdn_url).await { - Ok(true) => { - is_valid = true; - debug!( - "ImageCache: CDN URL verified successfully for {}", - original_url - ); - break; - } - Ok(false) => { - last_verify_error = "CDN returned non-image data or 404".to_string(); - } - Err(e) => { - last_verify_error = e; - } - } - - if attempt < 10 { - // Progressive backoff: 1s, 2s, 3s... up to 10s - let delay = 1000 * attempt; - tokio::time::sleep(std::time::Duration::from_millis(delay as u64)).await; - } - } - - if !is_valid { - error!( - "ImageCache: CDN verification failed for {} after 10 attempts: {}", - cdn_url, last_verify_error - ); - return Err(format!( - "CDN link was not accessible after upload: {}", - last_verify_error - )); - } - - // Save to database only after successful verification - self.repo.save_to_db(original_url, &cdn_url).await?; - - // Cache in Redis - let _ = self - .repo - .set_in_redis(&cache_key, &cdn_url, IMAGE_CACHE_TTL) - .await; - - // Invalidate API caches - let _ = self - .repo - .invalidate_api_caches(vec!["anime:*", "anime2:*", "komik:*"]) - .await; - - Ok(cdn_url) - } - .await; - - // Release Redis lock - let _ = self.repo.release_lock(&lock_key).await; - - work_result - } - .await; - - // Broadcast result - let _ = tx.send(result.clone()); - - // Remove from map - IN_FLIGHT_UPLOADS.remove(original_url); - - result - } - - /// Get CDN URL without uploading (read-only lookup) - pub async fn get_cdn_url(&self, original_url: &str) -> Option { - let cache_key = format!("{}:{}", IMAGE_CACHE_PREFIX, url_hash(original_url)); - - // Check Redis first - if let Some(cached_url) = self.repo.get_from_redis(&cache_key).await { - return Some(cached_url); - } - - // Check database - if let Ok(Some(cdn_url)) = self.repo.get_from_db(original_url).await { - return Some(cdn_url); - } - - None - } - - /// Find an original URL for a given CDN URL (reverse lookup) - pub async fn find_original_from_cdn(&self, cdn_url: &str) -> Option { - self.repo - .find_original_from_cdn(cdn_url) - .await - .ok() - .flatten() - } - - /// Invalidate cache for a URL - pub async fn invalidate(&self, original_url: &str) -> Result<(), String> { - let cache_key = format!("{}:{}", IMAGE_CACHE_PREFIX, url_hash(original_url)); - - // Remove from Redis - let _ = self.repo.delete_from_redis(&cache_key).await; - - // Remove from database - self.repo.delete_from_db(original_url).await?; - - debug!("ImageCache: Invalidated {}", original_url); - Ok(()) - } - - /// Helper to perform a single upload attempt - async fn perform_single_upload( - &self, - api_url: &str, - image_bytes: &[u8], - filename: &str, - ) -> Result { - debug!("ImageCache: Attempting upload to API server: {}", api_url); - - let part = reqwest::multipart::Part::bytes(image_bytes.to_vec()) - .file_name(filename.to_string()) - .mime_str("image/jpeg") - .map_err(|e| { - let err = format!("Failed to create multipart form for {}: {}", api_url, e); - error!("ImageCache: {}", err); - err - })?; - - let form = reqwest::multipart::Form::new().part("file", part); - - let response = self - .client - .post(api_url) - .multipart(form) - .send() - .await - .map_err(|e| { - let err = format!("Failed to send request to Picser API ({}): {}", api_url, e); - error!("ImageCache: {}", err); - err - })?; - - let response_status = response.status(); - let response_text = response.text().await.map_err(|e| { - let err = format!( - "Failed to read Picser response from {} (Status {}): {}", - api_url, response_status, e - ); - error!("ImageCache: {}", err); - err - })?; - - // Raw responses stay at debug level for troubleshooting without noisy default logs - debug!( - "ImageCache: Raw response from {} (Status {}): {}", - api_url, response_status, response_text - ); - - if !response_status.is_success() { - let error_message = serde_json::from_str::(&response_text) - .ok() - .and_then(|value| { - value - .get("error") - .and_then(|error| error.as_str()) - .map(|error| error.to_string()) - .or_else(|| { - value - .get("message") - .and_then(|message| message.as_str()) - .map(|message| message.to_string()) - }) - }) - .unwrap_or_else(|| response_text.clone()); - - let err = format!( - "Picser upload failed at {} (HTTP {}): {}", - api_url, response_status, error_message - ); - error!("ImageCache: {}", err); - return Err(err); - } - - let picser_response: PicserResponse = - serde_json::from_str(&response_text).map_err(|e| { - let err = format!( - "Failed to parse Picser response from {}: {} - Raw: {}", - api_url, e, response_text - ); - error!("ImageCache: {}", err); - err - })?; - - if !picser_response.success { - let err_msg = picser_response - .error - .unwrap_or_else(|| "Unknown error".to_string()); - let err = format!( - "Picser upload failed at {} (server error): {}", - api_url, err_msg - ); - error!("ImageCache: {}", err); - return Err(err); - } - - debug!("ImageCache: Upload successful to API server: {}", api_url); - Ok(picser_response) - } - - fn picser_api_endpoints(&self) -> Vec { - let configured = CONFIG.urls.picser_api_url.clone(); - let mut endpoints = vec![STATIC_PICSER_API_ENDPOINTS[0].to_string()]; - if !configured.is_empty() && !endpoints.contains(&configured) { - endpoints.push(configured); - } - for endpoint in STATIC_PICSER_API_ENDPOINTS.iter().skip(1) { - let endpoint = endpoint.to_string(); - if !endpoints.contains(&endpoint) { - endpoints.push(endpoint); - } - } - endpoints - } - - async fn upload_to_fallback_api( - &self, - image_bytes: &[u8], - filename: &str, - ) -> Result { - debug!( - "ImageCache: Attempting fallback upload to: {}", - CONFIG.urls.fallback_upload_api_url - ); - - let part = reqwest::multipart::Part::bytes(image_bytes.to_vec()) - .file_name(filename.to_string()) - .mime_str("image/jpeg") - .map_err(|e| { - let err = format!("Failed to create fallback multipart form: {}", e); - error!("ImageCache: {}", err); - err - })?; - - let form = reqwest::multipart::Form::new() - .part("file", part) - .text("fileName", filename.to_string()); - - let response = self - .client - .post(&CONFIG.urls.fallback_upload_api_url) - .multipart(form) - .send() - .await - .map_err(|e| { - let err = format!( - "Failed to send request to fallback upload API ({}): {}", - CONFIG.urls.fallback_upload_api_url, e - ); - error!("ImageCache: {}", err); - err - })?; - - let response_status = response.status(); - let response_text = response.text().await.map_err(|e| { - let err = format!( - "Failed to read fallback upload response from {} (Status {}): {}", - CONFIG.urls.fallback_upload_api_url, response_status, e - ); - error!("ImageCache: {}", err); - err - })?; - - debug!( - "ImageCache: Raw response from fallback upload API (Status {}): {}", - response_status, response_text - ); - - if !response_status.is_success() { - let err = format!( - "Fallback upload failed at {} (HTTP {}): {}", - CONFIG.urls.fallback_upload_api_url, response_status, response_text - ); - error!("ImageCache: {}", err); - return Err(err); - } - - let fallback_response: FallbackUploadResponse = serde_json::from_str(&response_text) - .map_err(|e| { - let err = format!( - "Failed to parse fallback upload response from {}: {} - Raw: {}", - CONFIG.urls.fallback_upload_api_url, e, response_text - ); - error!("ImageCache: {}", err); - err - })?; - - if fallback_response.download_url.trim().is_empty() { - return Err("Fallback upload response did not include download_url".to_string()); - } - - debug!( - "ImageCache: Fallback upload successful - URL: {}", - fallback_response.download_url - ); - Ok(fallback_response.download_url) - } - - async fn download_image_bytes(&self, original_url: &str) -> Result { - let mut candidates = vec![original_url.to_string()]; - if original_url.contains("https://alqanime.net/wp-content/") { - candidates.push(original_url.replace("https://alqanime.net/", "https://alqanime.si/")); - } - - let mut last_error = String::from("No download attempt started"); - for candidate in candidates { - debug!( - "ImageCache: Starting image download from source: {}", - candidate - ); - match self.client.get(&candidate).send().await { - Ok(response) => match response.bytes().await { - Ok(bytes) => { - let is_valid_image = infer::get(&bytes) - .map(|kind| kind.mime_type().starts_with("image/")) - .unwrap_or(false); - if is_valid_image { - return Ok(bytes); - } - - let trace_preview = String::from_utf8_lossy(&bytes) - .chars() - .take(100) - .collect::(); - last_error = format!( - "Image source ({}) returned non-image data (Preview: {})", - candidate, trace_preview - ); - warn!("ImageCache: {}", last_error); - } - Err(e) => { - last_error = format!( - "Failed to read image bytes from source ({}): {}", - candidate, e - ); - warn!("ImageCache: {}", last_error); - } - }, - Err(e) => { - last_error = format!( - "Failed to download image from source ({}): {}", - candidate, e - ); - warn!("ImageCache: {}", last_error); - } - } - } - - error!("ImageCache: {}", last_error); - Err(last_error) - } - - /// Upload image to Picser CDN with fallback upload API support - async fn upload_to_picser(&self, original_url: &str) -> Result { - // Download the image first - let image_bytes = self.download_image_bytes(original_url).await?; - - debug!( - "ImageCache: Image downloaded successfully, size: {} bytes", - image_bytes.len() - ); - - // Determine filename from URL - let filename = self.extract_filename(original_url); - - debug!( - "ImageCache: Will attempt upload to {} API endpoints sequentially with failover", - self.picser_api_endpoints().len() - ); - - let mut last_failed_api = String::from("Unknown"); - let picser_api_endpoints = self.picser_api_endpoints(); - let picser_api_count = picser_api_endpoints.len(); - - for (attempt_num, api_url) in picser_api_endpoints.iter().enumerate() { - let attempt_number = attempt_num + 1; - debug!( - "ImageCache: [Attempt {}/{}] Uploading {} bytes to: {}", - attempt_number, - picser_api_count, - image_bytes.len(), - api_url - ); - - match tokio::time::timeout( - std::time::Duration::from_secs(30), - self.perform_single_upload(api_url, &image_bytes, &filename), - ) - .await - { - Ok(Ok(response)) => match self.extract_cdn_url(response, original_url) { - Ok(cdn_url) => { - debug!( - "ImageCache: Upload succeeded on attempt {}/{} - CDN URL: {}", - attempt_number, picser_api_count, cdn_url - ); - return Ok(cdn_url); - } - Err(e) => { - last_failed_api = api_url.to_string(); - warn!( - "ImageCache: [Attempt {}/{}] Upload from {} did not yield a CDN URL: {}", - attempt_number, - picser_api_count, - api_url, - e - ); - error!( - "ImageCache: Attempt {}/{} failed - Last failed API: {} - Error: {}", - attempt_number, picser_api_count, last_failed_api, e - ); - } - }, - Ok(Err(e)) => { - last_failed_api = api_url.to_string(); - warn!( - "ImageCache: [Attempt {}/{}] Upload to {} failed: {}", - attempt_number, picser_api_count, api_url, e - ); - error!( - "ImageCache: Attempt {}/{} failed - Last failed API: {} - Error: {}", - attempt_number, picser_api_count, last_failed_api, e - ); - } - Err(_) => { - last_failed_api = api_url.to_string(); - let err = format!("Timeout (30s) while uploading to API endpoint: {}", api_url); - warn!( - "ImageCache: [Attempt {}/{}] {}", - attempt_number, picser_api_count, err - ); - error!( - "ImageCache: Attempt {}/{} failed - Last failed API: {} - Error: {}", - attempt_number, picser_api_count, last_failed_api, err - ); - } - } - } - - warn!( - "ImageCache: All {} Picser upload attempts failed for source URL: {}. Trying fallback upload API: {}", - picser_api_count, - original_url, - CONFIG.urls.fallback_upload_api_url - ); - - match tokio::time::timeout( - std::time::Duration::from_secs(30), - self.upload_to_fallback_api(&image_bytes, &filename), - ) - .await - { - Ok(Ok(cdn_url)) => Ok(cdn_url), - Ok(Err(e)) => { - error!( - "ImageCache: Fallback upload API failed after Picser failures. Last Picser API: {} - Fallback error: {}", - last_failed_api, e - ); - Err(format!( - "All {} Picser upload attempts and fallback upload API failed. Last Picser API endpoint: {}. Fallback error: {}", - picser_api_count, - last_failed_api, - e - )) - } - Err(_) => { - let err = format!( - "Timeout (30s) while uploading to fallback API endpoint: {}", - CONFIG.urls.fallback_upload_api_url - ); - error!( - "ImageCache: Fallback upload API timed out after Picser failures. Last Picser API: {} - {}", - last_failed_api, err - ); - Err(format!( - "All {} Picser upload attempts and fallback upload API failed. Last Picser API endpoint: {}. Fallback error: {}", - picser_api_count, - last_failed_api, - err - )) - } - } - } - - /// Internally verify a CDN URL's accessibility and validity - pub async fn verify_cdn_url(&self, cdn_url: &str) -> Result { - let resp = self.client.get(cdn_url).send().await.map_err(|e| { - let err = format!("Network error verifying CDN URL ({}): {}", cdn_url, e); - warn!("ImageCache: {}", err); - err - })?; - - let status = resp.status(); - if !status.is_success() { - let err = format!( - "CDN verification failed with HTTP {} for URL: {}", - status, cdn_url - ); - warn!("ImageCache: {}", err); - return Ok(false); - } - - let bytes = resp.bytes().await.map_err(|e| { - let err = format!("Failed to read bytes from CDN URL ({}): {}", cdn_url, e); - warn!("ImageCache: {}", err); - err - })?; - - // Structural verification (Fast MIME check) - let is_valid = infer::get(&bytes) - .map(|k| k.mime_type().starts_with("image/")) - .unwrap_or(false); - - if is_valid { - debug!( - "ImageCache: CDN URL verified successfully - content is valid image: {}", - cdn_url - ); - Ok(true) - } else { - warn!( - "ImageCache: CDN URL verification failed - content is not a valid image ({}): {}", - cdn_url, - String::from_utf8_lossy(&bytes[0..std::cmp::min(100, bytes.len())]) - ); - Ok(false) - } - } - - /// Extract CDN URL from Picser response - fn extract_cdn_url( - &self, - response: PicserResponse, - original_url: &str, - ) -> Result { - // Try to extract CDN URL from various response fields - if let Some(urls) = &response.urls { - if let Some(url) = &urls.raw { - debug!( - "ImageCache: Using CDN URL from urls.raw for source: {}", - original_url - ); - return Ok(url.clone()); - } - if let Some(url) = &urls.jsdelivr_commit { - debug!( - "ImageCache: Using CDN URL from urls.jsdelivr_commit for source: {}", - original_url - ); - return Ok(url.clone()); - } - if let Some(url) = &urls.jsdelivr { - debug!( - "ImageCache: Using CDN URL from urls.jsdelivr for source: {}", - original_url - ); - return Ok(url.clone()); - } - if let Some(url) = &urls.github { - debug!( - "ImageCache: Using CDN URL from urls.github for source: {}", - original_url - ); - return Ok(url.clone()); - } - } - - if let Some(url) = &response.url { - debug!( - "ImageCache: Using CDN URL from response.url for source: {}", - original_url - ); - return Ok(url.clone()); - } - - if let Some(url) = &response.github_url { - debug!( - "ImageCache: Using CDN URL from response.github_url for source: {}", - original_url - ); - return Ok(url.clone()); - } - - // If we get here, no CDN URL was found in the response - error!( - "ImageCache: No CDN URL found in Picser API response for source URL: {}. Response fields - success: {}, has_urls: {}, has_url: {}, has_github_url: {}, error: {}", - original_url, - response.success, - response.urls.is_some(), - response.url.is_some(), - response.github_url.is_some(), - response.error.as_deref().unwrap_or("none") - ); - Err(format!( - "No CDN URL in Picser response. Checked: urls.{{jsdelivr_commit,jsdelivr,raw,github}}, url, github_url. Response error: {}", - response.error.unwrap_or_else(|| "none".to_string()) - )) - } - - /// Extract filename from URL - fn extract_filename(&self, url: &str) -> String { - url.split('/') - .last() - .and_then(|s| s.split('?').next()) - .filter(|s| !s.is_empty() && s.contains('.')) - .map(|s| s.to_string()) - .unwrap_or_else(|| format!("{}.jpg", url_hash(url))) - } -} - -/// Convenience function to create a CDN URL for an image -/// Returns the original URL if caching fails (graceful fallback) -pub async fn cache_image_url( - db: Arc, - redis: &RedisPool, - original_url: &str, -) -> String { - let repo = Arc::new(SeaOrmImageCacheRepository::new(db, redis.clone())); - let cache = ImageCache::new(repo); - match cache.get_or_cache(original_url).await { - Ok(cdn_url) => cdn_url, - Err(e) => { - warn!("ImageCache: Failed to cache {}: {}", original_url, e); - to_wp_cdn(original_url) // Use WP CDN as graceful fallback - } - } -} - -/// Batch cache multiple images -pub async fn cache_image_urls( - db: Arc, - redis: &RedisPool, - urls: &[String], -) -> Vec { - let repo = Arc::new(SeaOrmImageCacheRepository::new(db, redis.clone())); - let cache = ImageCache::new(repo); - let mut results = Vec::with_capacity(urls.len()); - - for url in urls { - let cdn_url = match cache.get_or_cache(url).await { - Ok(u) => u, - Err(_) => url.clone(), - }; - results.push(cdn_url); - } - - results -} - -/// Helper to convert image URL to CDN URL in background (non-blocking) -/// Returns original URL immediately and caches in background -pub fn cache_image_url_lazy( - db: Arc, - redis: &RedisPool, - original_url: String, - semaphore: Option>, -) -> String { - let db_owned = db; - let redis_owned = redis.clone(); - let url = original_url.clone(); - let sem_owned = semaphore.clone(); - - // Spawn background task to cache - tokio::spawn(async move { - let repo = Arc::new(SeaOrmImageCacheRepository::new(db_owned, redis_owned)); - let mut cache = ImageCache::new(repo); - if let Some(sem) = sem_owned { - cache = cache.with_semaphore(sem); - } - - match cache.get_or_cache(&url).await { - Ok(_) => {} - Err(_) => {} - } - }); - - to_wp_cdn(&original_url) -} - -/// Convert image URL to CDN URL if already cached, otherwise return original -/// and trigger background caching for next request (with duplicate prevention) -/// Convert image URL to CDN URL if already cached, otherwise return original -/// and trigger background caching for next request (with duplicate prevention) -pub async fn get_cached_or_original( - db: Arc, - redis: &RedisPool, - original_url: &str, - semaphore: Option>, -) -> String { - let repo = Arc::new(SeaOrmImageCacheRepository::new(db.clone(), redis.clone())); - let cache = ImageCache::new(repo); - - // Check if already cached (Redis or DB) - if let Some(cdn_url) = cache.get_cdn_url(original_url).await { - return cdn_url; - } - - // Check if currently being cached by another process - let lock_key = format!("{}:{}", IMAGE_CACHE_LOCK_PREFIX, url_hash(original_url)); - let redis_cache = Cache::new(redis); - if redis_cache.get::(&lock_key).await.is_some() { - return to_wp_cdn(original_url); - } - - // Not cached and not being cached - start background caching - let db_owned = db.clone(); - let redis_owned = redis.clone(); - let url = original_url.to_string(); - let sem_owned = semaphore.clone(); - - tokio::spawn(async move { - let repo = Arc::new(SeaOrmImageCacheRepository::new(db_owned, redis_owned)); - let mut cache = ImageCache::new(repo); - if let Some(sem) = sem_owned { - cache = cache.with_semaphore(sem); - } - - let _ = cache.get_or_cache(&url).await; - }); - - to_wp_cdn(original_url) -} - -/// Batch process multiple image URLs - returns original URLs immediately -/// and triggers background caching for all -/// Batch process multiple image URLs - checks cache first, returns cached URL if found -/// For misses: returns original URL and triggers background caching -pub async fn cache_image_urls_batch_lazy( - db: Arc, - redis: &RedisPool, - urls: Vec, - semaphore: Option>, -) -> Vec { - if urls.is_empty() { - return urls; - } - - let mut results = vec![String::new(); urls.len()]; - let mut missing_indices = Vec::new(); - - // 1. Batch check Redis - let redis_cache = Cache::new(redis); - let cache_keys: Vec = urls - .iter() - .map(|url| format!("{}:{}", IMAGE_CACHE_PREFIX, url_hash(url))) - .collect(); - - let cached_values: Vec> = redis_cache.mget(&cache_keys).await; - - for (i, val) in cached_values.iter().enumerate() { - if let Some(cdn_url) = val { - results[i] = cdn_url.clone(); - } else { - missing_indices.push(i); - } - } - - // 2. Batch check Database for Redis misses - if !missing_indices.is_empty() { - let missing_urls: Vec = missing_indices.iter().map(|&i| urls[i].clone()).collect(); - - // Note: Using repository for batch check would be better, but keeping it direct for now to match SeaORM usage - // but I should probably add a batch method to repository later. - use crate::infrastructure::persistence::entities::image_cache; - use sea_orm::{ColumnTrait, EntityTrait, QueryFilter}; - - match image_cache::Entity::find() - .filter(image_cache::Column::OriginalUrl.is_in(missing_urls.clone())) - .all(db.as_ref()) - .await - { - Ok(db_entries) => { - let db_map: std::collections::HashMap = db_entries - .into_iter() - .map(|e| (e.original_url, e.cdn_url)) - .collect(); - - let mut still_missing_indices = Vec::new(); - - for &idx in &missing_indices { - let url = &urls[idx]; - if let Some(cdn_url) = db_map.get(url) { - results[idx] = cdn_url.clone(); - // Put back to Redis - let _ = redis_cache - .set_with_ttl(&cache_keys[idx], cdn_url, IMAGE_CACHE_TTL) - .await; - } else { - // Real miss - return WP CDN proxy and trigger background upload - results[idx] = to_wp_cdn(url); - still_missing_indices.push(idx); - } - } - - // 3. Trigger background caching for still missing URLs - if !still_missing_indices.is_empty() { - let db_owned = db.clone(); - let redis_owned = redis.clone(); - let sem_owned = semaphore.clone(); - let urls_to_cache: Vec = still_missing_indices - .iter() - .map(|&idx| urls[idx].clone()) - .collect(); - - tokio::spawn(async move { - use futures::stream::{self, StreamExt}; - let repo = Arc::new(SeaOrmImageCacheRepository::new(db_owned, redis_owned)); - let mut cache = ImageCache::new(repo); - if let Some(sem) = sem_owned { - cache = cache.with_semaphore(sem); - } - - stream::iter(urls_to_cache) - .map(|url| { - let cache_ref = &cache; - async move { - let _ = cache_ref.get_or_cache(&url).await; - } - }) - .buffer_unordered(20) - .collect::>() - .await; - }); - } - } - Err(e) => { - error!("ImageCache: Batch DB check failed: {}", e); - for &idx in &missing_indices { - results[idx] = to_wp_cdn(&urls[idx]); - } - } - } - } - - results -} - -/// Apply cached CDN poster URLs to a collection of items using the HasPoster trait. -pub async fn apply_cached_posters( - items: &mut [T], - db: Arc, - redis: &RedisPool, - semaphore: Option>, -) { - let posters: Vec = items.iter().map(|item| item.poster().to_string()).collect(); - let cached = cache_image_urls_batch_lazy(db, redis, posters, semaphore).await; - for (i, item) in items.iter_mut().enumerate() { - if let Some(url) = cached.get(i) { - item.set_poster(url.clone()); - } - } -} diff --git a/src/infrastructure/services/images/mod.rs b/src/infrastructure/services/images/mod.rs deleted file mode 100644 index a5c08fd..0000000 --- a/src/infrastructure/services/images/mod.rs +++ /dev/null @@ -1 +0,0 @@ -pub mod cache; diff --git a/src/infrastructure/services/mod.rs b/src/infrastructure/services/mod.rs deleted file mode 100644 index 8f0da9f..0000000 --- a/src/infrastructure/services/mod.rs +++ /dev/null @@ -1 +0,0 @@ -pub mod images; diff --git a/src/lib.rs b/src/lib.rs index ac00d26..2a65617 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -27,4 +27,3 @@ pub mod bootstrap; pub mod config; pub mod events; pub mod observability; -pub mod scheduler; diff --git a/src/observability/openapi_modules.rs b/src/observability/openapi_modules.rs index 08e771f..ce678e1 100644 --- a/src/observability/openapi_modules.rs +++ b/src/observability/openapi_modules.rs @@ -41,8 +41,6 @@ use utoipa::OpenApi; crate::presentation::handler::komik::search_slug, crate::presentation::handler::komik::search_slug_page, // Proxy module handlers - crate::presentation::handler::proxy::fetch_with_proxy_only, - crate::presentation::handler::proxy::image_cache, ), components( schemas( diff --git a/src/presentation/handler/anime.rs b/src/presentation/handler/anime.rs index 80fe6f7..1f74792 100644 --- a/src/presentation/handler/anime.rs +++ b/src/presentation/handler/anime.rs @@ -78,12 +78,7 @@ pub struct GenreListResponse { // ============================================================================ fn make_use_cases(state: &Arc) -> AnimeUseCases { - AnimeUseCases::new( - OtakudesuRepository::new(), - state.redis_pool.clone(), - state.db.clone(), - Some(state.image_processing_semaphore.clone()), - ) + AnimeUseCases::new(OtakudesuRepository::new(), state.redis_pool.clone()) } // ============================================================================ diff --git a/src/presentation/handler/anime2.rs b/src/presentation/handler/anime2.rs index 7b7ca69..60778cb 100644 --- a/src/presentation/handler/anime2.rs +++ b/src/presentation/handler/anime2.rs @@ -39,12 +39,7 @@ pub struct FilterQuery { // ============================================================================ fn make_use_cases(state: &Arc) -> Anime2UseCases { - Anime2UseCases::new( - AlqanimeRepository::new(), - state.redis_pool.clone(), - state.db.clone(), - Some(state.image_processing_semaphore.clone()), - ) + Anime2UseCases::new(AlqanimeRepository::new(), state.redis_pool.clone()) } // ============================================================================ diff --git a/src/presentation/handler/komik.rs b/src/presentation/handler/komik.rs index b945230..1f58b3b 100644 --- a/src/presentation/handler/komik.rs +++ b/src/presentation/handler/komik.rs @@ -25,12 +25,7 @@ use crate::presentation::state::AppState; // ============================================================================ fn make_use_cases(state: &Arc) -> KomikUseCases { - KomikUseCases::new( - KomikRepository::new(), - state.redis_pool.clone(), - state.db.clone(), - Some(state.image_processing_semaphore.clone()), - ) + KomikUseCases::new(KomikRepository::new(), state.redis_pool.clone()) } // ============================================================================ diff --git a/src/presentation/handler/proxy.rs b/src/presentation/handler/proxy.rs index 8a3f38a..392fdd6 100644 --- a/src/presentation/handler/proxy.rs +++ b/src/presentation/handler/proxy.rs @@ -1,131 +1 @@ -//! Proxy and image cache API handlers. - -use std::sync::Arc; - -use axum::extract::{Json, Query, State}; -use axum::response::Response; -use serde::Deserialize; -use tracing::info; -use utoipa::{IntoParams, ToSchema}; - -use crate::application::proxy::use_cases::{ - AuditImageCacheResult, ImageCacheResult, ProxyUseCases, -}; -use crate::infrastructure::repository::image_cache_seaorm::SeaOrmImageCacheRepository; -use crate::infrastructure::repository::ProxyRepository; -use crate::presentation::error::AppError; -use crate::presentation::state::AppState; - -// ============================================================================ -// Request DTOs -// ============================================================================ - -/// Query parameters for proxy fetch (GET). -#[derive(Debug, Deserialize, IntoParams, ToSchema)] -pub struct ProxyParams { - /// URL to fetch via proxy. - pub url: String, -} - -/// Request body for image cache (POST). -#[derive(Debug, Deserialize, ToSchema)] -pub struct ImageCacheRequest { - /// Original image URL to cache. - pub url: String, - /// If true, returns original URL immediately and caches in background. - #[serde(default)] - pub lazy: bool, -} - -/// Request body for auditing image cache (POST). -#[derive(Debug, Deserialize, ToSchema)] -pub struct AuditImageCacheRequest { - /// Original image URL to audit. - pub url: String, -} - -// ============================================================================ -// Helper -// ============================================================================ - -fn make_use_cases(state: &Arc) -> ProxyUseCases { - let repo = Arc::new(SeaOrmImageCacheRepository::new( - state.db.clone(), - state.redis_pool.clone(), - )); - ProxyUseCases::new(ProxyRepository::new(), repo) -} - -// ============================================================================ -// Handlers -// ============================================================================ - -/// GET /api/proxy/croxy โ€” Fetch a URL through the proxy and return raw bytes. -#[utoipa::path( - get, - path = "/api/proxy/croxy", - tag = "proxy", - operation_id = "proxy_croxy", - params(ProxyParams), - responses( - (status = 200, description = "Proxied response", body = Vec::, content_type = "application/octet-stream"), - (status = 500, description = "Internal Server Error"), - ) -)] -pub async fn fetch_with_proxy_only( - State(state): State>, - Query(params): Query, -) -> Result { - info!("Handling proxy fetch for URL: {}", params.url); - let response = make_use_cases(&state) - .fetch_with_proxy_only(params.url) - .await?; - Ok(response) -} - -/// POST /api/proxy/image-cache โ€” Cache an image URL to CDN. -#[utoipa::path( - post, - path = "/api/proxy/image-cache", - tag = "proxy", - operation_id = "proxy_image_cache", - request_body = ImageCacheRequest, - responses( - (status = 200, description = "Image cache result", body = ImageCacheResult), - (status = 500, description = "Internal Server Error"), - ) -)] -pub async fn image_cache( - State(state): State>, - Json(req): Json, -) -> Result, AppError> { - info!( - "Handling image cache for URL: {} (lazy: {})", - req.url, req.lazy - ); - let result = make_use_cases(&state) - .image_cache(req.url, req.lazy) - .await?; - Ok(Json(result)) -} - -/// POST /api/proxy/image-cache/audit โ€” Audit and repair a cached image. -#[utoipa::path( - post, - path = "/api/proxy/image-cache/audit", - tag = "proxy", - operation_id = "proxy_image_cache_audit", - request_body = AuditImageCacheRequest, - responses( - (status = 200, description = "Audit result", body = AuditImageCacheResult), - (status = 500, description = "Internal Server Error"), - ) -)] -pub async fn audit_image_cache( - State(state): State>, - Json(req): Json, -) -> Result, AppError> { - info!("Handling audit cache for URL: {}", req.url); - let result = make_use_cases(&state).audit_image_cache(req.url).await?; - Ok(Json(result)) -} +//! Proxy handlers. diff --git a/src/presentation/router.rs b/src/presentation/router.rs index 3e4116c..995758e 100644 --- a/src/presentation/router.rs +++ b/src/presentation/router.rs @@ -154,18 +154,7 @@ pub fn build_router(app_state: Arc) -> anyhow::Result { axum::routing::get(crate::presentation::handler::komik::search_slug_page), ) // Proxy routes - .route( - "/api/proxy/croxy", - axum::routing::get(crate::presentation::handler::proxy::fetch_with_proxy_only), - ) - .route( - "/api/proxy/image-cache", - axum::routing::post(crate::presentation::handler::proxy::image_cache), - ) - .route( - "/api/proxy/image-cache/audit", - axum::routing::post(crate::presentation::handler::proxy::audit_image_cache), - ) + // (all proxy routes removed) // Health .route( "/health", diff --git a/src/presentation/state.rs b/src/presentation/state.rs index 345ebca..7b1bb74 100644 --- a/src/presentation/state.rs +++ b/src/presentation/state.rs @@ -6,7 +6,6 @@ use deadpool_redis::Pool; use sea_orm::DatabaseConnection; use crate::events::bus::EventBus; -use crate::infrastructure::repository::SeaOrmImageCacheRepository; /// Shared application state injected into every handler via Axum State. /// @@ -16,9 +15,7 @@ use crate::infrastructure::repository::SeaOrmImageCacheRepository; pub struct AppState { pub redis_pool: Pool, pub db: Arc, - pub image_processing_semaphore: Arc, pub event_bus: Arc, - pub image_cache_repo: Arc, } impl AppState { diff --git a/src/scheduler/cleanup_cache.rs b/src/scheduler/cleanup_cache.rs deleted file mode 100644 index d753a5f..0000000 --- a/src/scheduler/cleanup_cache.rs +++ /dev/null @@ -1,222 +0,0 @@ -//! Scheduled task for cleaning up old cached data. - -use async_trait::async_trait; -use sea_orm::*; -use std::sync::Arc; -use tracing::{info, warn}; - -use crate::infrastructure::cache::redis::Cache; -use crate::infrastructure::cache::redis_pool::get_redis_pool; -use crate::infrastructure::persistence::entities::image_cache; - -use super::ScheduledTask; - -/// Cleanup old cache data to prevent disk/memory bloat. -/// Runs daily at 2 AM to clean: -/// - Old image cache entries (>30 days) -/// - Orphaned cache keys in Redis -/// - Expired data without TTL -pub struct CleanupOldCache { - db: Arc, -} - -impl CleanupOldCache { - pub fn new(db: Arc) -> Self { - Self { db } - } -} - -#[async_trait] -impl ScheduledTask for CleanupOldCache { - fn name(&self) -> &'static str { - "cleanup_old_cache" - } - - fn schedule(&self) -> &'static str { - // Daily at 2 AM - "0 0 2 * * *" - } - - async fn run(&self) { - tracing::debug!("๐Ÿงน Starting old cache cleanup..."); - - let mut total_cleaned = 0; - - // 1. Clean old image cache (>30 days) - match self.cleanup_old_images(30).await { - Ok(count) => { - if count > 0 { - info!("โœ“ Cleaned {} old image cache entries", count); - } - total_cleaned += count; - } - Err(e) => { - warn!("Failed to clean old image cache: {}", e); - } - } - - // 2. Clean orphaned Redis keys - match self.cleanup_orphaned_redis_keys().await { - Ok(count) => { - if count > 0 { - info!("โœ“ Cleaned {} orphaned Redis keys", count); - } - total_cleaned += count; - } - Err(e) => { - warn!("Failed to clean orphaned Redis keys: {}", e); - } - } - - // 3. Compact Redis memory - match self.compact_redis_memory().await { - Ok(()) => { - if total_cleaned > 0 { - info!("โœ“ Redis memory compacted"); - } - } - Err(e) => { - warn!("Failed to compact Redis memory: {}", e); - } - } - - if total_cleaned > 0 { - info!("๐ŸŽ‰ Cache cleanup complete: {} items cleaned", total_cleaned); - } - } -} - -impl CleanupOldCache { - /// Remove image cache entries older than specified days. - async fn cleanup_old_images(&self, days: i64) -> Result { - use chrono::{Duration, Utc}; - - let cutoff = Utc::now() - Duration::days(days); - - // Find old entries - let old_images = image_cache::Entity::find() - .filter(image_cache::Column::CreatedAt.lt(cutoff)) - .all(self.db.as_ref()) - .await - .map_err(|e| e.to_string())?; - - let count = old_images.len(); - - if count == 0 { - return Ok(0); - } - - // Delete from database - let ids: Vec = old_images.iter().map(|img| img.id.clone()).collect(); - - image_cache::Entity::delete_many() - .filter(image_cache::Column::Id.is_in(ids)) - .exec(self.db.as_ref()) - .await - .map_err(|e| e.to_string())?; - - // Also clean from Redis - let redis_pool = get_redis_pool().map_err(|e| e.to_string())?; - let cache = Cache::new(redis_pool); - for img in old_images { - let cache_key = format!("img_cache:{}", Self::hash_url(&img.original_url)); - let _ = cache.delete(&cache_key).await; - } - - Ok(count) - } - - /// Clean orphaned Redis keys (keys without TTL that shouldn't exist). - async fn cleanup_orphaned_redis_keys(&self) -> Result { - use deadpool_redis::redis::AsyncCommands; - - let pool = get_redis_pool().map_err(|e| e.to_string())?; - let mut conn = pool - .get() - .await - .map_err(|e| format!("Failed to get Redis connection: {}", e))?; - - let mut cleaned = 0; - - // Find keys without TTL (should not exist) - let patterns = vec!["anime:*", "komik:*", "user:*:profile", "img_cache:*"]; - - for pattern in patterns { - let keys = { - let mut iter: deadpool_redis::redis::AsyncIter<'_, String> = conn - .scan_match(pattern) - .await - .map_err(|e| format!("Failed to scan keys: {}", e))?; - - let mut keys = Vec::new(); - while let Some(key_result) = iter.next_item().await { - if let Ok(key) = key_result { - keys.push(key); - } - } - - keys - }; - - for key in keys { - let ttl: i64 = conn.ttl(key.as_str()).await.unwrap_or(-1); - - // TTL = -1 means no expiration (orphaned) - // TTL = -2 means key doesn't exist - if ttl == -1 { - // Set a default TTL of 7 days for orphaned keys - let _: () = conn.expire(key.as_str(), 604800).await.unwrap_or(()); - cleaned += 1; - } - } - } - - Ok(cleaned) - } - - /// Compact Redis memory to free up fragmented space. - async fn compact_redis_memory(&self) -> Result<(), String> { - let pool = get_redis_pool().map_err(|e| e.to_string())?; - let mut conn = pool - .get() - .await - .map_err(|e| format!("Failed to get Redis connection: {}", e))?; - - // Run MEMORY PURGE command - using cmd method on connection - let _: String = deadpool_redis::redis::cmd("MEMORY") - .arg("PURGE") - .query_async(&mut *conn) - .await - .map_err(|e| format!("Failed to purge memory: {}", e))?; - - Ok(()) - } - - /// Simple hash function for URL (same as in image_cache.rs). - /// Simple hash function for URL (same as in image_cache.rs). - fn hash_url(url: &str) -> String { - use sha2::{Digest, Sha256}; - let mut hasher = Sha256::new(); - hasher.update(url.as_bytes()); - format!("{:x}", hasher.finalize()) - } -} - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn test_schedule() { - use sea_orm::{DatabaseBackend, MockDatabase}; - let db = MockDatabase::new(DatabaseBackend::Sqlite).into_connection(); - let task = CleanupOldCache { db: Arc::new(db) }; - assert_eq!(task.schedule(), "0 0 2 * * *"); - } - - #[test] - fn test_hash_url() { - let hash = CleanupOldCache::hash_url("https://example.com/image.jpg"); - assert_eq!(hash.len(), 64); // SHA256 hash length - } -} diff --git a/src/scheduler/mod.rs b/src/scheduler/mod.rs deleted file mode 100644 index 7387dd2..0000000 --- a/src/scheduler/mod.rs +++ /dev/null @@ -1,5 +0,0 @@ -pub mod cleanup_cache; -pub mod runner; - -pub use cleanup_cache::CleanupOldCache; -pub use runner::{ScheduledTask, Scheduler}; diff --git a/src/scheduler/runner.rs b/src/scheduler/runner.rs deleted file mode 100644 index 13be76c..0000000 --- a/src/scheduler/runner.rs +++ /dev/null @@ -1,93 +0,0 @@ -//! Scheduler implementation using tokio-cron-scheduler. - -use async_trait::async_trait; -use std::sync::Arc; -use tokio_cron_scheduler::{Job, JobScheduler}; -use tracing::info; - -/// Trait for scheduled tasks. -#[async_trait] -pub trait ScheduledTask: Send + Sync { - /// Task name for logging. - fn name(&self) -> &'static str; - - /// Cron expression (e.g., "0 * * * * *" for every minute). - fn schedule(&self) -> &'static str; - - /// Execute the task. - async fn run(&self); -} - -/// Scheduler for running cron jobs. -pub struct Scheduler { - inner: JobScheduler, -} - -impl Scheduler { - /// Create a new scheduler. - pub async fn new() -> anyhow::Result { - let scheduler = JobScheduler::new().await?; - Ok(Self { inner: scheduler }) - } - - /// Add a task to the scheduler. - pub async fn add(&self, task: T) -> anyhow::Result<()> { - let task = Arc::new(task); - let task_name = task.name(); - let schedule = task.schedule(); - - let job = Job::new_async(schedule, move |_uuid, _lock| { - let task = Arc::clone(&task); - Box::pin(async move { - tracing::debug!("Running scheduled task: {}", task.name()); - task.run().await; - }) - })?; - - self.inner.add(job).await?; - info!("Scheduled task '{}' with cron: {}", task_name, schedule); - Ok(()) - } - - /// Add a simple job with a closure. - pub async fn add_job( - &self, - name: &'static str, - schedule: &str, - f: F, - ) -> anyhow::Result<()> - where - F: Fn() -> Fut + Send + Sync + 'static, - Fut: std::future::Future + Send + 'static, - { - let f = Arc::new(f); - let job = Job::new_async(schedule, move |_uuid, _lock| { - let f = Arc::clone(&f); - Box::pin(async move { - info!("Running scheduled job: {}", name); - f().await; - }) - })?; - - self.inner.add(job).await?; - info!("Scheduled job '{}' with cron: {}", name, schedule); - Ok(()) - } - - /// Start the scheduler. - pub async fn start(&self) -> anyhow::Result<()> { - info!("Starting scheduler"); - self.inner.start().await?; - Ok(()) - } - - /// Stop the scheduler. - pub async fn shutdown(&mut self) -> anyhow::Result<()> { - info!("Shutting down scheduler"); - self.inner.shutdown().await?; - Ok(()) - } -} - -// Real scheduled tasks with actual implementations -// Real scheduled tasks with actual implementations