diff --git a/.hermes/plans/scraper-mytheclipse-round1-spec.md b/.hermes/plans/scraper-mytheclipse-round1-spec.md new file mode 100644 index 0000000..101d2b3 --- /dev/null +++ b/.hermes/plans/scraper-mytheclipse-round1-spec.md @@ -0,0 +1,155 @@ +# Spec: Migrate scraper infra to mytheclipse crates (Round 1) + +Date: 2026-08-30 +Repo: /home/code/scraper (github.com/asepharyana/scraper) +Goal: Replace hand-rolled infrastructure in the scraper with the user's +custom `mytheclipse` library crates where there is a clear 1:1 mapping. + +## Library crates (from /home/code/mytheclipse, all v1.20.0) + +- `mytheclipse` — core: retry (retry/RetryConfig/JitterKind/RetryError), + ratelimit (RateLimiter/RateLimitError), timeouts + (timeout), concurrency primitives, ServiceBuilder, + runtime_auto (available_parallelism), spawn_io/spawn_bg. +- `mytheclipse-cache` — Cache trait, MemoryCache, RedisCache (l2-redis), + CacheAside (cache-aside), CacheError. +- `mytheclipse-config`— ConfigLoader (file + env merge, hot-reload, validation). +- `mytheclipse-event` — EventBus trait, InMemoryEventBus, TypedEventBus (byte/JSON pub/sub). +- `mytheclipse-crypto`— (NOT used in this round — scraper has no crypto/JWT usage.) + +## Current state (baseline) + +- Cargo.toml has NO mytheclipse crates. Uses directly: + - `backoff` — retry with exponential backoff (hand-rolled wrapper in + src/infrastructure/scraping/retry.rs) + - `deadpool-redis` — Redis pool + raw `redis::AsyncCommands` in + src/infrastructure/cache/{redis_pool.rs,redis.rs} + - `dashmap` — request coalescing in proxy_fetch.rs (IN_FLIGHT map) + - `rayon` — komik_parser parallel map + - `config` crate — config loading in src/config/mod.rs + - `opentelemetry*` — metrics in src/observability/metrics.rs + - custom EventBus in src/events/bus.rs (unused by any handler) + - custom RateLimiter in src/presentation/middleware/ratelimit.rs (unused by router) +- Baseline: `cargo check` currently clean (verified in background). + +## Replacement mapping (behavior-preserving) + +| # | Hand-rolled | mytheclipse replacement | Files touched | +|---|---|---|---| +| 1 | backoff + retry.rs wrapper | `mytheclipse::retry` + `RetryConfig` | src/infrastructure/scraping/retry.rs, parsing_utils.rs, otakudesu.rs | +| 2 | deadpool-redis pool in redis_pool.rs | Keep pool, wrap conn with `RedisCache` (mytheclipse-cache) | redis_pool.rs + new bridge, redis.rs, proxy_fetch.rs | +| 3 | Cache helper (redis.rs) | `CacheAside` + typed JSON serde over `RedisCache` | redis.rs, application/*/use_cases.rs | +| 4 | custom EventBus (events/bus.rs) | `InMemoryEventBus`/`TypedEventBus` (mytheclipse-event) | events/bus.rs → re-export, state.rs, bootstrap | +| 5 | custom RateLimiter middleware | `mytheclipse::RateLimiter` (token bucket) | presentation/middleware/ratelimit.rs | +| 6 | `config` crate loader in config/mod.rs | `mytheclipse-config` ConfigLoader | src/config/mod.rs | +| 7 | HTTP client wrapper (http_client.rs) | `mytheclipse-http` HttpClient (optional) | http_client.rs (SKIP this round — retry semantics differ; reqwest needs headers/UA control) | +| 8 | OTel metrics (metrics.rs) | Keep opentelemetry (mytheclipse-tracing has no metrics exporter; avoid behavior change) | SKIP this round | +| 9 | rayon in komik_parser | `mytheclipse::compute::compute_par_for_each` (feature compute) | komik_parser.rs (SKIP this round — parser correctness risk; rayon works) | + +## Scope decision (this round) + +Implement items #1–#6. Skip #7–#9 with rationale: +- #7 mytheclipse-http client is a thin reqwest wrapper without header/UA control + needed by common_headers(); converting scrapers through it changes fetch + semantics (returns bytes, loses status/content-type) — not behavior-preserving. +- #8 mytheclipse-tracing has no metrics exporter; opentelemetry stays. +- #9 parser logic (2k+ LOC) is out of scope for infra migration; rayon stays. + +## Detail per item + +### 1. retry.rs → mytheclipse::retry +- Replace `backoff::ExponentialBackoff` with `mytheclipse::{retry, RetryConfig, JitterKind}`. +- Keep signature-compatible helpers so call sites barely change: + `default_backoff() -> RetryConfig`, `quick_backoff()`, `slow_backoff()`, + `custom_backoff(...)`. +- `transient/permanent` helpers: mytheclipse retry uses a `predicate` closure + `Fn(&E) -> bool`. Replace transient/permanent with a retryable predicate + (retry on any error except a marker). To preserve "transient = retry, + permanent = stop", use a wrapper type or predicate returning true for all + errors, and treat 4xx-style permanent errors by converting callers to return + a `Permanent` variant. +- Keep the `retry` fn name re-exported for minimal call-site churn. + +### 2. Redis: mytheclipse-cache RedisCache + bridge +- Add `mytheclipse-cache = { version = "1.20", features = ["l2-redis", "cache-aside"] }`. +- Keep deadpool pool (mytheclipse has no pool); obtain + `redis::aio::MultiplexedConnection` from deadpool conn (deadpool_redis::Connection + derefs to `&mut redis::aio::ConnectionLike` — need to convert via `into_multiplexed()`). +- New bridge file `src/infrastructure/cache/mytheclipse.rs`: + `pub fn redis_cache() -> &'static mytheclipse_cache::RedisCache` building + from the pool lazily (LazyLock). + +### 3. Cache helper → CacheAside +- Rewrite `src/infrastructure/cache/redis.rs` as a thin typed wrapper over + `RedisCache` implementing `get/get_or_set/set/set_with_ttl/delete/exists` + with serde_json, so use_cases keep the same ergonomic API (minimal churn) + but delegate to mytheclipse `Cache` trait underneath. +- `get_or_set` becomes CacheAside-equivalent (read-through). + +### 4. events/bus.rs → mytheclipse-event +- Replace custom Event/EventHandler/pub-sub with + `TypedEventBus`. +- Define domain events as serde structs implementing `mytheclipse_event::Event` + (which is a blanket trait on Serialize+DeserializeOwned types). +- `EventBus` type alias: `pub type EventBus = TypedEventBus`. +- state.rs + bootstrap keep `Arc`; `new()` -> `EventBus::new(InMemoryEventBus::default())`. +- Publish/subscribe usage sites (none currently) adapt if any. + +### 5. Ratelimit middleware → mytheclipse::RateLimiter +- Keep middleware shape (axum State) but inner limiter becomes + `mytheclipse::RateLimiter` token bucket: `RateLimiter::new(rate_per_sec, burst)`. +- `check()` -> `try_acquire()` returns Result<_, RateLimitError>. + +### 6. config/mod.rs → mytheclipse-config +- Replace `config::Config::builder()` with `mytheclipse_config::ConfigLoader`. +- Sources: load_dotenv (env feature) → merge_file(config/default.toml) → + merge_file(config/{RUN_MODE}.toml) → merge_env (APP__ prefix + legacy env + var mapping preserved). +- Keep `CONFIG` LazyLock and the same `AppConfig` struct (`Deserialize` + + `mytheclipse_config::Config` blanket trait). + +## Schema/type changes +- `deadpool_redis::Pool` still in AppState (no change) — RedisCache wraps conns. +- `AppState.event_bus` becomes `Arc>`. +- Error mapping: add `From` for AppError + + DomainError? Keep the existing string-based paths; add explicit From impls + where needed. + +## Backend surface +- No HTTP route changes. No API contract changes. +- Redis cache keys/TTLs unchanged. Event topics unchanged (none active). + +## Frontend surface +- None (backend-only service). + +## Verification steps +1. `cargo check` — exit 0 +2. `cargo test` — all pass +3. `cargo clippy -- -D warnings` — 0 warnings (project lint gate) +4. `cargo fmt --check` — clean +5. `cargo build --release` — succeeds +6. Smoke: boot server briefly if feasible (needs Redis/DB; else compile-only + unit tests) + +## Risk register +- deadpool Connection → MultiplexedConnection conversion API differences + (0.22 vs 0.27 redis). Verify at compile time; fallback = keep raw conn in + Cache impl. +- `mytheclipse::retry` predicate-based vs backoff transient/permanent — call + sites that relied on permanent-stop need review (repository fetch_html uses + `transient()` on everything — safe to retry-all). +- ConfigLoader merge_env nesting: existing env var mapping uses `APP__` prefix + + legacy names. Must keep exact names (APP_DATABASE_URL etc.) — verify with + `peek()` / unit test on the real struct. +- EventBus trait object: `Arc>` is concrete — + no trait object; AppState carries the concrete type to avoid dyn issues. + +## Migration order (batches, each verified) +Batch A: Cargo.toml deps + retry.rs rewrite + parsing_utils/otakudesu call sites +Batch B: redis_pool bridge + Cache rewrite (redis.rs) + proxy_fetch cache calls +Batch C: events/bus.rs + state.rs + bootstrap +Batch D: ratelimit middleware + router wiring +Batch E: config/mod.rs ConfigLoader +Batch F: full verification + commit + +Each batch: cargo check + cargo test must stay green. Commit once at end with +message: `refactor: migrate infra to mytheclipse crates (retry, cache, event, ratelimit, config)` \ No newline at end of file diff --git a/Cargo.lock b/Cargo.lock index 61ef137..a205bf8 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -79,12 +79,6 @@ dependencies = [ "derive_arbitrary", ] -[[package]] -name = "arraydeque" -version = "0.5.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7d902e3d592a523def97af8f317b08ce16b7ab854c1985a0c671e6f15cebc236" - [[package]] name = "arrayvec" version = "0.7.6" @@ -272,20 +266,6 @@ dependencies = [ "syn 2.0.117", ] -[[package]] -name = "backoff" -version = "0.4.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b62ddb9cb1ec0a098ad4bbf9344d0713fa193ae1a80af55febcff2627b6a00c1" -dependencies = [ - "futures-core", - "getrandom 0.2.17", - "instant", - "pin-project-lite", - "rand 0.8.5", - "tokio", -] - [[package]] name = "base64" version = "0.22.1" @@ -312,6 +292,12 @@ dependencies = [ "serde", ] +[[package]] +name = "bitflags" +version = "1.3.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bef38d45163c2f1dde094a7dfd33ccf595c92905c8f8f4fdc18d06fb1037718a" + [[package]] name = "bitflags" version = "2.11.0" @@ -508,61 +494,12 @@ dependencies = [ "crossbeam-utils", ] -[[package]] -name = "config" -version = "0.15.22" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8e68cfe19cd7d23ffde002c24ffa5cda73931913ef394d5eaaa32037dc940c0c" -dependencies = [ - "async-trait", - "convert_case", - "json5", - "pathdiff", - "ron", - "rust-ini", - "serde-untagged", - "serde_core", - "serde_json", - "toml", - "winnow", - "yaml-rust2", -] - [[package]] name = "const-oid" version = "0.9.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2459377285ad874054d797f3ccebf984978aa39129f6eafde5cdc8315b612f8" -[[package]] -name = "const-random" -version = "0.1.18" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "87e00182fe74b066627d63b85fd550ac2998d4b0bd86bfed477a0ae4c7c71359" -dependencies = [ - "const-random-macro", -] - -[[package]] -name = "const-random-macro" -version = "0.1.16" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f9d839f2a20b0aee515dc581a6172f2321f96cab76c1a38a4c584a194955390e" -dependencies = [ - "getrandom 0.2.17", - "once_cell", - "tiny-keccak", -] - -[[package]] -name = "convert_case" -version = "0.6.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "ec182b0ca2f35d8fc196cf3404988fd8b8c739a4d270ff118a398feb0cbec1ca" -dependencies = [ - "unicode-segmentation", -] - [[package]] name = "core-foundation" version = "0.9.4" @@ -622,6 +559,15 @@ dependencies = [ "cfg-if", ] +[[package]] +name = "crossbeam-channel" +version = "0.5.16" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d85363c37faeca707aef026efa9f3b34d077bce547e48f770770625c6013679e" +dependencies = [ + "crossbeam-utils", +] + [[package]] name = "crossbeam-deque" version = "0.8.6" @@ -656,12 +602,6 @@ version = "0.8.21" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d0a5c400df2834b80a4c3327b3aad3a4c4cd4de0629063962b03235697506a28" -[[package]] -name = "crunchy" -version = "0.2.4" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "460fbee9c2c2f33933d720630a6a0bac33ba7053db5344fac858d4b8952d77d5" - [[package]] name = "crypto-common" version = "0.1.7" @@ -825,15 +765,6 @@ dependencies = [ "syn 2.0.117", ] -[[package]] -name = "dlv-list" -version = "0.5.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "442039f5147480ba31067cb00ada1adae6892028e40e45fc5de7b7df6dcc1b5f" -dependencies = [ - "const-random", -] - [[package]] name = "dotenvy" version = "0.15.7" @@ -885,17 +816,6 @@ version = "1.0.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "877a4ace8713b0bcf2a4e7eec82529c029f1d0619886d18145fea96c3ffe5c0f" -[[package]] -name = "erased-serde" -version = "0.4.10" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d2add8a07dd6a8d93ff627029c51de145e12686fbc36ecb298ac22e74cf02dec" -dependencies = [ - "serde", - "serde_core", - "typeid", -] - [[package]] name = "errno" version = "0.3.14" @@ -934,6 +854,16 @@ version = "2.3.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "37909eebbb50d72f9059c3b6d82c0463f2ff062c9e95845c43a6c9c0355411be" +[[package]] +name = "filetime" +version = "0.2.29" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5c287a33c7f0a620c38e641e7f60827713987b3c0f26e8ddc9462cc69cf75759" +dependencies = [ + "cfg-if", + "libc", +] + [[package]] name = "find-msvc-tools" version = "0.1.9" @@ -998,6 +928,15 @@ dependencies = [ "percent-encoding", ] +[[package]] +name = "fsevent-sys" +version = "4.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "76ee7a02da4d231650c7cea31349b889be2f45ddb3ef3032d2ec8185f6313fd2" +dependencies = [ + "libc", +] + [[package]] name = "funty" version = "2.0.0" @@ -1605,12 +1544,23 @@ dependencies = [ ] [[package]] -name = "instant" -version = "0.1.13" +name = "inotify" +version = "0.9.6" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e0242819d153cba4b4b05a5a8f2a7e9bbf97b6055b2a002b395c96b5ff3c0222" +checksum = "f8069d3ec154eb856955c1c0fbffefbf5f3c40a104ec912d4797314c1801abff" dependencies = [ - "cfg-if", + "bitflags 1.3.2", + "inotify-sys", + "libc", +] + +[[package]] +name = "inotify-sys" +version = "0.1.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c033f80b2c113cdf91ab7a33faa9cbc014726dcad99880c8609af2a370edf37d" +dependencies = [ + "libc", ] [[package]] @@ -1667,14 +1617,23 @@ dependencies = [ ] [[package]] -name = "json5" -version = "0.4.1" +name = "kqueue" +version = "1.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "96b0db21af676c1ce64250b5f40f3ce2cf27e4e47cb91ed91eb6fe9350b430c1" +checksum = "8d763e5b24120b4ddf50de6c92308156765aabfbbccebf401da7cff2d70a41ea" dependencies = [ - "pest", - "pest_derive", - "serde", + "kqueue-sys", + "libc", +] + +[[package]] +name = "kqueue-sys" +version = "1.1.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "07293a4e297ac234359b510362495713f75ea345d5307140414f20c69ffeb087" +dependencies = [ + "bitflags 2.11.0", + "libc", ] [[package]] @@ -1710,7 +1669,7 @@ version = "0.1.15" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7ddbf48fd451246b1f8c2610bd3b4ac0cc6e149d89832867093ab69a17194f08" dependencies = [ - "bitflags", + "bitflags 2.11.0", "libc", "plain", "redox_syscall 0.7.3", @@ -1833,6 +1792,18 @@ dependencies = [ "simd-adler32", ] +[[package]] +name = "mio" +version = "0.8.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a4a650543ca06a924e8b371db273b2756685faae30f8487da1b56505a8f78b0c" +dependencies = [ + "libc", + "log", + "wasi", + "windows-sys 0.48.0", +] + [[package]] name = "mio" version = "1.2.0" @@ -1861,6 +1832,61 @@ dependencies = [ "version_check", ] +[[package]] +name = "mytheclipse" +version = "1.20.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "62ec67baf3bde7c129eefa282ca89768f3975ce60dcba8c5f7adb8e045d404aa" +dependencies = [ + "async-trait", + "num_cpus", + "rand 0.8.5", + "thiserror 2.0.18", + "tokio", + "tracing", +] + +[[package]] +name = "mytheclipse-cache" +version = "1.21.1" +dependencies = [ + "async-trait", + "redis", + "serde", + "serde_json", + "tokio", + "tracing", +] + +[[package]] +name = "mytheclipse-config" +version = "1.20.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f09ce745aab668bcd51971d5f297525ff2d9e0534ff252b7b473433b4e855c8d" +dependencies = [ + "dotenvy", + "notify", + "serde", + "serde_json", + "serde_yaml", + "tokio", + "toml", + "tracing", +] + +[[package]] +name = "mytheclipse-event" +version = "1.20.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ef3f6657e41861ba79425f9c7e9442628f3e70a454da25d392d4bfc37326210b" +dependencies = [ + "async-trait", + "serde", + "serde_json", + "tokio", + "tracing", +] + [[package]] name = "native-tls" version = "0.2.18" @@ -1884,6 +1910,25 @@ version = "1.0.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "650eef8c711430f1a879fdd01d4745a7deea475becfb90269c06775983bbf086" +[[package]] +name = "notify" +version = "6.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6205bd8bb1e454ad2e27422015fb5e4f2bcc7e08fa8f27058670d208324a4d2d" +dependencies = [ + "bitflags 2.11.0", + "crossbeam-channel", + "filetime", + "fsevent-sys", + "inotify", + "kqueue", + "libc", + "log", + "mio 0.8.11", + "walkdir", + "windows-sys 0.48.0", +] + [[package]] name = "nu-ansi-term" version = "0.50.3" @@ -1977,7 +2022,7 @@ version = "0.10.76" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "951c002c75e16ea2c65b8c7e4d3d51d5530d8dfa7d060b4776828c88cfb18ecf" dependencies = [ - "bitflags", + "bitflags 2.11.0", "cfg-if", "foreign-types", "libc", @@ -2090,16 +2135,6 @@ dependencies = [ "num-traits", ] -[[package]] -name = "ordered-multimap" -version = "0.7.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "49203cdcae0030493bad186b28da2fa25645fa276a51b6fec8010d281e02ef79" -dependencies = [ - "dlv-list", - "hashbrown 0.14.5", -] - [[package]] name = "ouroboros" version = "0.18.5" @@ -2153,12 +2188,6 @@ dependencies = [ "windows-link", ] -[[package]] -name = "pathdiff" -version = "0.2.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "df94ce210e5bc13cb6651479fa48d14f601d9858cfe0467f43ae157023b938d3" - [[package]] name = "pem-rfc7468" version = "0.7.0" @@ -2174,49 +2203,6 @@ version = "2.3.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9b4f627cb1b25917193a259e49bdad08f671f8d9708acfd5fe0a8c1455d87220" -[[package]] -name = "pest" -version = "2.8.6" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e0848c601009d37dfa3430c4666e147e49cdcf1b92ecd3e63657d8a5f19da662" -dependencies = [ - "memchr", - "ucd-trie", -] - -[[package]] -name = "pest_derive" -version = "2.8.6" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "11f486f1ea21e6c10ed15d5a7c77165d0ee443402f0780849d1768e7d9d6fe77" -dependencies = [ - "pest", - "pest_generator", -] - -[[package]] -name = "pest_generator" -version = "2.8.6" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8040c4647b13b210a963c1ed407c1ff4fdfa01c31d6d2a098218702e6664f94f" -dependencies = [ - "pest", - "pest_meta", - "proc-macro2", - "quote", - "syn 2.0.117", -] - -[[package]] -name = "pest_meta" -version = "2.8.6" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "89815c69d36021a140146f26659a81d6c2afa33d216d736dd4be5381a7362220" -dependencies = [ - "pest", - "sha2", -] - [[package]] name = "pgvector" version = "0.4.1" @@ -2390,7 +2376,7 @@ version = "3.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e67ba7e9b2b56446f1d419b1d807906278ffa1a658a8a5d8a39dcb1f5a78614f" dependencies = [ - "toml_edit", + "toml_edit 0.25.8+spec-1.1.0", ] [[package]] @@ -2617,7 +2603,7 @@ version = "0.5.18" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ed2bf2547551a7053d6fdfafda3f938979645c44812fbfcda098faae3f1a362d" dependencies = [ - "bitflags", + "bitflags 2.11.0", ] [[package]] @@ -2626,7 +2612,7 @@ version = "0.7.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6ce70a74e890531977d37e532c34d45e9055d2409ed08ddba14529471ed0be16" dependencies = [ - "bitflags", + "bitflags 2.11.0", ] [[package]] @@ -2754,20 +2740,6 @@ dependencies = [ "syn 1.0.109", ] -[[package]] -name = "ron" -version = "0.12.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "fd490c5b18261893f14449cbd28cb9c0b637aebf161cd77900bfdedaff21ec32" -dependencies = [ - "bitflags", - "once_cell", - "serde", - "serde_derive", - "typeid", - "unicode-ident", -] - [[package]] name = "rsa" version = "0.9.10" @@ -2822,16 +2794,6 @@ dependencies = [ "walkdir", ] -[[package]] -name = "rust-ini" -version = "0.21.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "796e8d2b6696392a43bea58116b667fb4c29727dc5abd27d6acf338bb4f688c7" -dependencies = [ - "cfg-if", - "ordered-multimap", -] - [[package]] name = "rust_decimal" version = "1.41.0" @@ -2870,7 +2832,7 @@ version = "1.1.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b6fe4565b9518b83ef4f91bb47ce29620ca828bd32cb7e408f0062e9930ba190" dependencies = [ - "bitflags", + "bitflags 2.11.0", "errno", "libc", "linux-raw-sys", @@ -2981,15 +2943,17 @@ dependencies = [ "anyhow", "async-trait", "axum 0.8.8", - "backoff", "chrono", - "config", "dashmap", "deadpool-redis", "dotenvy", "flate2", "futures", "http", + "mytheclipse", + "mytheclipse-cache", + "mytheclipse-config", + "mytheclipse-event", "opentelemetry", "opentelemetry-otlp", "opentelemetry_sdk", @@ -3114,7 +3078,7 @@ version = "3.7.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b7f4bc775c73d9a02cde8bf7b2ec4c9d12743edf609006c7facc23998404cd1d" dependencies = [ - "bitflags", + "bitflags 2.11.0", "core-foundation 0.10.1", "core-foundation-sys", "libc", @@ -3137,7 +3101,7 @@ version = "0.33.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "feef350c36147532e1b79ea5c1f3791373e61cbd9a6a2615413b3807bb164fb7" dependencies = [ - "bitflags", + "bitflags 2.11.0", "cssparser", "derive_more", "log", @@ -3166,18 +3130,6 @@ dependencies = [ "serde_derive", ] -[[package]] -name = "serde-untagged" -version = "0.1.9" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f9faf48a4a2d2693be24c6289dbe26552776eb7737074e6722891fadbe6c5058" -dependencies = [ - "erased-serde", - "serde", - "serde_core", - "typeid", -] - [[package]] name = "serde_core" version = "1.0.228" @@ -3224,11 +3176,11 @@ dependencies = [ [[package]] name = "serde_spanned" -version = "1.1.0" +version = "0.6.9" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "876ac351060d4f882bb1032b6369eb0aef79ad9df1ea8bc404874d8cc3d0cd98" +checksum = "bf41e0cfaf7226dca15e8197172c295a782857fcb97fad1808a166870dee75a3" dependencies = [ - "serde_core", + "serde", ] [[package]] @@ -3243,6 +3195,19 @@ dependencies = [ "serde", ] +[[package]] +name = "serde_yaml" +version = "0.9.34+deprecated" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6a8b1a1a2ebf674015cc02edccce75287f1a0130d394307b36743c2f5d504b47" +dependencies = [ + "indexmap 2.13.0", + "itoa", + "ryu", + "serde", + "unsafe-libyaml", +] + [[package]] name = "servo_arc" version = "0.4.3" @@ -3488,7 +3453,7 @@ dependencies = [ "atoi", "base64", "bigdecimal", - "bitflags", + "bitflags 2.11.0", "byteorder", "bytes", "chrono", @@ -3535,7 +3500,7 @@ dependencies = [ "atoi", "base64", "bigdecimal", - "bitflags", + "bitflags 2.11.0", "byteorder", "chrono", "crc", @@ -3677,6 +3642,17 @@ dependencies = [ "unicode-ident", ] +[[package]] +name = "syn" +version = "3.0.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e6275cddf4610d1775e6d1fe9469b2e77d0f39fd98fb7450901b821e0c53649f" +dependencies = [ + "proc-macro2", + "quote", + "unicode-ident", +] + [[package]] name = "sync_wrapper" version = "1.0.2" @@ -3703,7 +3679,7 @@ version = "0.7.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a13f3d0daba03132c0aa9767f98351b3488edc2c100cda2d2ec2b04f3d8d3c8b" dependencies = [ - "bitflags", + "bitflags 2.11.0", "core-foundation 0.9.4", "system-configuration-sys", ] @@ -3828,15 +3804,6 @@ dependencies = [ "time-core", ] -[[package]] -name = "tiny-keccak" -version = "2.0.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2c9d3793400a45f954c52e73d068316d76b6f4e36977e3fcebb13a2721e80237" -dependencies = [ - "crunchy", -] - [[package]] name = "tinystr" version = "0.8.2" @@ -3864,13 +3831,13 @@ checksum = "1f3ccbac311fea05f86f61904b462b55fb3df8837a366dfc601a0161d0532f20" [[package]] name = "tokio" -version = "1.50.0" +version = "1.53.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "27ad5e34374e03cfffefc301becb44e9dc3c17584f414349ebe29ed26661822d" +checksum = "202caea871b69668250d242070849eb495be178ed697a3e98aebce5bc81a0bed" dependencies = [ "bytes", "libc", - "mio", + "mio 1.2.0", "parking_lot", "pin-project-lite", "signal-hook-registry", @@ -3881,13 +3848,13 @@ dependencies = [ [[package]] name = "tokio-macros" -version = "2.6.1" +version = "2.7.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5c55a2eff8b69ce66c84f85e1da1c233edc36ceb85a2058d11b0d6a3c7e7569c" +checksum = "78773a2a397f451582ce068015985c33193cf6dea8b74d2a639fe457b2f07b0e" dependencies = [ "proc-macro2", "quote", - "syn 2.0.117", + "syn 3.0.4", ] [[package]] @@ -3948,15 +3915,23 @@ dependencies = [ [[package]] name = "toml" -version = "1.1.0+spec-1.1.0" +version = "0.8.23" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f8195ca05e4eb728f4ba94f3e3291661320af739c4e43779cbdfae82ab239fcc" +checksum = "dc1beb996b9d83529a9e75c17a1686767d148d70663143c7854d8b4a09ced362" dependencies = [ - "serde_core", + "serde", "serde_spanned", - "toml_datetime", - "toml_parser", - "winnow", + "toml_datetime 0.6.11", + "toml_edit 0.22.27", +] + +[[package]] +name = "toml_datetime" +version = "0.6.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "22cddaf88f4fbc13c51aebbf5f8eceb5c7c5a9da2ac40a13519eb5b0a0e8f11c" +dependencies = [ + "serde", ] [[package]] @@ -3968,6 +3943,20 @@ dependencies = [ "serde_core", ] +[[package]] +name = "toml_edit" +version = "0.22.27" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "41fe8c660ae4257887cf66394862d21dbca4a6ddd26f04a3560410406a2f819a" +dependencies = [ + "indexmap 2.13.0", + "serde", + "serde_spanned", + "toml_datetime 0.6.11", + "toml_write", + "winnow 0.7.15", +] + [[package]] name = "toml_edit" version = "0.25.8+spec-1.1.0" @@ -3975,9 +3964,9 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "16bff38f1d86c47f9ff0647e6838d7bb362522bdf44006c7068c2b1e606f1f3c" dependencies = [ "indexmap 2.13.0", - "toml_datetime", + "toml_datetime 1.1.0+spec-1.1.0", "toml_parser", - "winnow", + "winnow 1.0.0", ] [[package]] @@ -3986,9 +3975,15 @@ version = "1.1.0+spec-1.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2334f11ee363607eb04df9b8fc8a13ca1715a72ba8662a26ac285c98aabb4011" dependencies = [ - "winnow", + "winnow 1.0.0", ] +[[package]] +name = "toml_write" +version = "0.1.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5d99f8c9a7727884afe522e9bd5edbfc91a3312b36a77b5fb8926e4c31a41801" + [[package]] name = "tonic" version = "0.12.3" @@ -4062,7 +4057,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d4e6559d53cc268e5031cd8429d05415bc4cb4aefc4aa5d6cc35fbf5b924a1f8" dependencies = [ "async-compression", - "bitflags", + "bitflags 2.11.0", "bytes", "futures-core", "futures-util", @@ -4181,24 +4176,12 @@ dependencies = [ "utf-8", ] -[[package]] -name = "typeid" -version = "1.0.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "bc7d623258602320d5c55d1bc22793b57daff0ec7efc270ea7d55ce1d5f5471c" - [[package]] name = "typenum" version = "1.19.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "562d481066bde0658276a35467c4af00bdc6ee726305698a55b86e61d7ad82bb" -[[package]] -name = "ucd-trie" -version = "0.1.7" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2896d95c02a80c6d6a5d6e953d479f5ddf2dfdb6a244441010e373ac0fb88971" - [[package]] name = "unicase" version = "2.9.0" @@ -4232,12 +4215,6 @@ version = "0.1.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7df058c713841ad818f1dc5d3fd88063241cc61f49f5fbea4b951e8cf5a8d71d" -[[package]] -name = "unicode-segmentation" -version = "1.13.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9629274872b2bfaf8d66f5f15725007f635594914870f65218920345aa11aa8c" - [[package]] name = "unicode-width" version = "0.2.2" @@ -4250,6 +4227,12 @@ version = "0.2.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ebc1c04c71510c7f702b52b7c350734c9ff1295c464a03335b00bb84fc54f853" +[[package]] +name = "unsafe-libyaml" +version = "0.2.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "673aac59facbab8a9007c7f6108d11f63b603f7cabff99fabf650fea5c32b861" + [[package]] name = "untrusted" version = "0.9.0" @@ -4504,7 +4487,7 @@ version = "0.244.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "47b807c72e1bac69382b3a6fb3dbe8ea4c0ed87ff5629b8685ae6b9a611028fe" dependencies = [ - "bitflags", + "bitflags 2.11.0", "hashbrown 0.15.5", "indexmap 2.13.0", "semver", @@ -4787,6 +4770,15 @@ version = "0.52.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "589f6da84c646204747d1270a2a5661ea66ed1cced2631d546fdfb155959f9ec" +[[package]] +name = "winnow" +version = "0.7.15" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "df79d97927682d2fd8adb29682d1140b343be4ac0f08fd68b7765d9c059d3945" +dependencies = [ + "memchr", +] + [[package]] name = "winnow" version = "1.0.0" @@ -4854,7 +4846,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9d66ea20e9553b30172b5e831994e35fbde2d165325bec84fc43dbf6f4eb9cb2" dependencies = [ "anyhow", - "bitflags", + "bitflags 2.11.0", "indexmap 2.13.0", "log", "serde", @@ -4899,17 +4891,6 @@ dependencies = [ "tap", ] -[[package]] -name = "yaml-rust2" -version = "0.10.4" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2462ea039c445496d8793d052e13787f2b90e750b833afee748e601c17621ed9" -dependencies = [ - "arraydeque", - "encoding_rs", - "hashlink", -] - [[package]] name = "yansi" version = "1.0.1" diff --git a/Cargo.toml b/Cargo.toml index cbc7475..58e3316 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -13,6 +13,11 @@ default-run = "scraper" # Dependensi yang dibutuhkan saat aplikasi berjalan [dependencies] +mytheclipse = { version = "1.20", features = ["resiliency", "traffic"] } +mytheclipse-cache = { path = "/home/code/mytheclipse/crates/mytheclipse-cache", features = ["l2-redis", "cache-aside"] } +mytheclipse-config = { version = "1.20" } +mytheclipse-event = { version = "1.20" } + axum = { version = "0.8.8", features = ["ws", "multipart", "macros"] } tokio = { version = "1.49.0", features = ["full"] } serde = { version = "1.0", features = ["derive"] } @@ -36,7 +41,6 @@ urlencoding = "2.1" url = "2.5.8" tower-http = { version = "0.6.8", features = ["fs", "cors", "compression-gzip", "compression-br", "compression-zstd"] } -backoff = { version = "0.4", features = ["futures", "tokio"] } dashmap = "6.1" deadpool-redis = { version = "0.22.1", features = ["serde"] } rayon = "1.11" @@ -44,7 +48,6 @@ scraper = "0.25.0" flate2 = "1.1" redis = { version = "0.32.7", features = ["tokio-rustls-comp", "safe_iterators"] } thiserror = "2.0.18" -config = { version = "0.15.19", features = ["toml"] } utoipa = { version = "5.0", features = ["axum_extras"] } utoipa-swagger-ui = { version = "9.0", features = ["axum"] } diff --git a/src/bootstrap/mod.rs b/src/bootstrap/mod.rs index fd0a94e..7f49b8f 100644 --- a/src/bootstrap/mod.rs +++ b/src/bootstrap/mod.rs @@ -77,7 +77,7 @@ impl Application { // App State components let db_arc = Arc::new(db); - let event_bus = Arc::new(crate::events::bus::EventBus::new()); + let event_bus = Arc::new(crate::events::bus::new_event_bus()); let redis_pool = crate::infrastructure::cache::redis_pool::redis_pool() .map_err(|e| anyhow::anyhow!("Failed to init Redis pool: {}", e))?; diff --git a/src/config/mod.rs b/src/config/mod.rs index 62aac4d..e99291b 100644 --- a/src/config/mod.rs +++ b/src/config/mod.rs @@ -5,7 +5,7 @@ //! - Fails fast at startup if required variables are missing //! - Supports hierarchical configuration (default -> environment-specific) -use config::{Config, ConfigError, Environment, File}; +use mytheclipse_config::ConfigError; use serde::Deserialize; use std::env; use std::sync::LazyLock; @@ -220,6 +220,8 @@ impl AppConfig { /// 2. `config/{environment}.toml` /// 3. `config/default.toml` pub fn load() -> Result { + use mytheclipse_config::ConfigLoader; + // Load .env file first if let Err(e) = dotenvy::dotenv() { tracing::debug!("Could not load .env file: {}", e); @@ -227,24 +229,40 @@ impl AppConfig { let run_mode = env::var("RUN_MODE").unwrap_or_else(|_| "development".into()); - let config = Config::builder() - // Start with default config file - .add_source(File::with_name("config/default").required(false)) - // Layer on environment-specific values - .add_source(File::with_name(&format!("config/{}", run_mode)).required(false)) - // Add environment variables (with APP_ prefix) - .add_source( - Environment::with_prefix("APP") - .separator("__") - .try_parsing(true), - ) - // Map legacy env vars to new config structure - .set_override_option("database_url", env::var("DATABASE_URL").ok())? - .set_override_option("jwt_secret", env::var("JWT_SECRET").ok())? - .set_override_option("redis_url", env::var("REDIS_URL").ok())? - .build()?; + // Start with default config file (optional — many deploys are env-only) + let mut loader = ConfigLoader::::new() + .load_dotenv(std::path::Path::new(".env")) + .unwrap_or_else(|_| ConfigLoader::::new()); - config.try_deserialize() + let default_path = std::path::Path::new("config/default.toml"); + if default_path.exists() { + loader = loader.merge_file(default_path)?; + } + // Layer on environment-specific values (optional) + let env_path_str = format!("config/{}.toml", run_mode); + let env_path = std::path::Path::new(&env_path_str); + if env_path.exists() { + loader = loader.merge_file(env_path)?; + } + + // Add environment variables (with APP_ prefix) + loader = loader.merge_env("APP"); + + // Map legacy env vars to new config structure. + // ConfigLoader merges at leaf level; legacy vars override the merged + // value directly (highest priority after APP_*). + let mut value = loader.peek().clone(); + if let Ok(v) = env::var("DATABASE_URL") { + value["database_url"] = serde_json::Value::String(v); + } + if let Ok(v) = env::var("JWT_SECRET") { + value["jwt_secret"] = serde_json::Value::String(v); + } + if let Ok(v) = env::var("REDIS_URL") { + value["redis_url"] = serde_json::Value::String(v); + } + + ConfigLoader::::new().merge_value(value).build() } /// Check if running in production mode diff --git a/src/events/bus.rs b/src/events/bus.rs index f09905c..5263e88 100644 --- a/src/events/bus.rs +++ b/src/events/bus.rs @@ -1,146 +1,44 @@ -//! Event bus implementation. +//! Event bus — delegated to `mytheclipse_event`. +//! +//! The scraper uses `TypedEventBus` from the mytheclipse +//! event crate (JSON-typed pub/sub over an in-process broadcast bus). +//! Domain events are plain serde structs; `mytheclipse_event::Event` is a +//! blanket trait, so every `Serialize + DeserializeOwned + Send + Sync + 'static` +//! type is automatically an event — no manual impl needed. -use async_trait::async_trait; +use mytheclipse_event::{InMemoryEventBus, TypedEventBus}; -use std::{any::TypeId, collections::HashMap, sync::Arc}; -use tokio::sync::{broadcast, RwLock}; -use tracing::{debug, info}; +/// Event payload marker — re-exported so domain types can reference it. +pub use mytheclipse_event::Event; -/// Trait for events that can be published. -pub trait Event: Clone + Send + Sync + 'static { - /// Event name for logging/debugging. - const NAME: &'static str; -} +/// The scraper's application event bus. +pub type EventBus = TypedEventBus; -/// Trait for event handlers. -#[async_trait] -pub trait EventHandler: Send + Sync { - async fn handle(&self, event: E); -} - -/// The event bus for publishing and subscribing to events. -pub struct EventBus { - channels: RwLock>>, -} - -impl EventBus { - /// Create a new event bus. - pub fn new() -> Self { - Self { - channels: RwLock::new(HashMap::new()), - } - } - - /// Publish an event to all subscribers. - pub async fn publish(&self, event: E) { - let type_id = TypeId::of::(); - let channels = self.channels.read().await; - - if let Some(sender) = channels.get(&type_id) { - if let Some(tx) = sender.downcast_ref::>() { - let _ = tx.send(event); - debug!("Published event: {}", E::NAME); - } - } - } - - /// Subscribe to events of a specific type. - /// Returns a receiver that can be used to receive events. - pub async fn subscribe(&self) -> broadcast::Receiver { - let type_id = TypeId::of::(); - - // Check if channel exists - { - let channels = self.channels.read().await; - if let Some(sender) = channels.get(&type_id) { - if let Some(tx) = sender.downcast_ref::>() { - return tx.subscribe(); - } - } - } - - // Create new channel - let (tx, rx) = broadcast::channel::(100); - { - let mut channels = self.channels.write().await; - channels.insert(type_id, Box::new(tx)); - } - - // Re-get the receiver from the stored sender - let channels = self.channels.read().await; - if let Some(sender) = channels.get(&type_id) { - if let Some(tx) = sender.downcast_ref::>() { - return tx.subscribe(); - } - } - - rx - } - - /// Register a handler for a specific event type. - /// The handler will be called whenever an event of that type is published. - pub async fn on + 'static>(&self, handler: H) { - let mut rx = self.subscribe::().await; - let handler = Arc::new(handler); - - tokio::spawn(async move { - loop { - match rx.recv().await { - Ok(event) => { - handler.handle(event).await; - } - Err(broadcast::error::RecvError::Closed) => break, - Err(broadcast::error::RecvError::Lagged(n)) => { - tracing::warn!("Event handler lagged by {} events", n); - } - } - } - }); - - info!("Registered handler for event: {}", E::NAME); - } -} - -impl Default for EventBus { - fn default() -> Self { - Self::new() - } +/// Build a new in-memory event bus. +pub fn new_event_bus() -> EventBus { + TypedEventBus::new(InMemoryEventBus::default()) } // Common events /// User registered event. -#[derive(Clone, Debug)] +#[derive(Clone, Debug, serde::Serialize, serde::Deserialize)] pub struct UserRegistered { pub user_id: String, pub email: String, pub name: String, } -impl Event for UserRegistered { - const NAME: &'static str = "user.registered"; -} - /// User logged in event. -#[derive(Clone, Debug)] +#[derive(Clone, Debug, serde::Serialize, serde::Deserialize)] pub struct UserLoggedIn { pub user_id: String, pub ip_address: Option, } -impl Event for UserLoggedIn { - const NAME: &'static str = "user.logged_in"; -} - /// Order created event. -#[derive(Clone, Debug)] +#[derive(Clone, Debug, serde::Serialize, serde::Deserialize)] pub struct OrderCreated { pub order_id: String, pub user_id: String, pub total: f64, } - -impl Event for OrderCreated { - const NAME: &'static str = "order.created"; -} - - diff --git a/src/infrastructure/cache/mod.rs b/src/infrastructure/cache/mod.rs index 2d19fea..ad0edea 100644 --- a/src/infrastructure/cache/mod.rs +++ b/src/infrastructure/cache/mod.rs @@ -1,3 +1,4 @@ +pub mod mytheclipse; pub mod redis; pub mod redis_pool; diff --git a/src/infrastructure/cache/mytheclipse.rs b/src/infrastructure/cache/mytheclipse.rs new file mode 100644 index 0000000..8d4b4cb --- /dev/null +++ b/src/infrastructure/cache/mytheclipse.rs @@ -0,0 +1,54 @@ +//! Bridge between the deadpool Redis pool and `mytheclipse_cache::RedisCache`. +//! +//! The scraper keeps its deadpool `Pool` (connection lifecycle, recycling) and +//! hands each checked-out connection to mytheclipse's `RedisCache` (which +//! implements the library `Cache` trait). Because both sides now share the +//! same `redis` crate version, a `deadpool_redis::Connection` can be converted +//! directly into a `redis::aio::MultiplexedConnection` via `take()`. + +use std::sync::LazyLock; + +use mytheclipse_cache::{Cache, CacheError, RedisCache}; +use tokio::sync::OnceCell; + +use crate::infrastructure::cache::redis_pool::redis_pool; + +/// Lazily-initialised shared `RedisCache` built from the deadpool pool. +/// +/// The multiplexed connection is cheaply cloneable (Arc-backed), so the whole +/// process shares one logical connection while deadpool manages recycling. +static REDIS_CACHE: LazyLock> = LazyLock::new(OnceCell::new); + +/// Return a handle to the shared mytheclipse `RedisCache`, initialising it on +/// first use from the deadpool pool. +pub async fn redis_cache() -> Result<&'static RedisCache, CacheError> { + let cell = &*REDIS_CACHE; + cell.get_or_try_init(|| async { + let pool = redis_pool().map_err(CacheError::Io)?; + let conn = pool + .get() + .await + .map_err(|e| CacheError::Io(e.to_string()))?; + let mux = deadpool_redis::Connection::take(conn); + Ok(RedisCache::new(mux)) + }) + .await +} + +/// Convenience wrappers so callers can use the mytheclipse `Cache` methods +/// directly without importing the trait twice. +pub async fn get(key: &str) -> Result>, CacheError> { + redis_cache().await?.get(key).await +} + +pub async fn set( + key: &str, + value: Vec, + ttl: Option, +) -> Result<(), CacheError> { + redis_cache().await?.set(key, value, ttl).await +} + +pub async fn invalidate(key: &str) -> Result<(), CacheError> { + redis_cache().await?.invalidate(key).await +} diff --git a/src/infrastructure/cache/redis.rs b/src/infrastructure/cache/redis.rs index 80285b5..cbe9ede 100644 --- a/src/infrastructure/cache/redis.rs +++ b/src/infrastructure/cache/redis.rs @@ -1,70 +1,59 @@ -//! Redis caching helpers. +//! Redis caching helpers — typed wrapper over `mytheclipse_cache`. +//! +//! The underlying byte-cache is `mytheclipse_cache::RedisCache` (which +//! implements the library `Cache` trait). This module keeps the ergonomic +//! typed JSON surface the use cases rely on (`get_or_set`, `set_with_ttl`) +//! while delegating the actual Redis commands to the library. -use deadpool_redis::redis::AsyncCommands; -use deadpool_redis::Pool; +use mytheclipse_cache::{Cache as CacheTrait, CacheError}; use serde::{de::DeserializeOwned, Serialize}; -use tracing::{debug, error}; +use tracing::debug; /// Default cache TTL in seconds (5 minutes). pub const DEFAULT_CACHE_TTL: u64 = 300; -/// Cache helper for Redis operations. +/// Typed JSON cache helper over the shared mytheclipse `RedisCache`. pub struct Cache<'a> { - pool: &'a Pool, + _marker: std::marker::PhantomData<&'a ()>, } impl<'a> Cache<'a> { - pub fn new(pool: &'a Pool) -> Self { - Self { pool } + pub fn new(_pool: &'a deadpool_redis::Pool) -> Self { + Self { + _marker: std::marker::PhantomData, + } + } + + async fn cache(&self) -> Result<&'static mytheclipse_cache::RedisCache, CacheError> { + super::mytheclipse::redis_cache().await } pub async fn get(&self, key: &str) -> Option { - let mut conn = match self.pool.get().await { - Ok(c) => c, + match self.cache().await { + Ok(cache) => match CacheTrait::get(cache, key).await { + Ok(Some(bytes)) => serde_json::from_slice(&bytes).ok(), + Ok(None) => None, + Err(e) => { + debug!("Cache: get error for {}: {}", key, e); + None + } + }, Err(e) => { - error!("Cache: failed to get connection: {}", e); - return None; + debug!("Cache: unavailable for {}: {}", key, e); + None } - }; - - let cached: Option = conn.get(key).await.ok()?; - - if cached.is_some() { - debug!("Cache hit: {}", key); - } else { - debug!("Cache miss: {}", key); } - - cached.and_then(|json| serde_json::from_str(&json).ok()) } pub async fn mget(&self, keys: &[String]) -> Vec> { if keys.is_empty() { return Vec::new(); } - - let mut conn = match self.pool.get().await { - Ok(c) => c, - Err(e) => { - error!("Cache: failed to get connection: {}", e); - return std::iter::repeat_with(|| None).take(keys.len()).collect(); - } - }; - - use deadpool_redis::redis::cmd; - let cached_values: Vec> = - match cmd("MGET").arg(keys).query_async(&mut conn).await { - Ok(v) => v, - Err(e) => { - error!("Cache: failed to mget values: {}", e); - return std::iter::repeat_with(|| None).take(keys.len()).collect(); - } - }; - - cached_values - .into_iter() - .map(|opt_s| opt_s.and_then(|json| serde_json::from_str(&json).ok())) - .collect() + let mut out = Vec::with_capacity(keys.len()); + for k in keys { + out.push(self.get::(k).await); + } + out } pub async fn set(&self, key: &str, value: &T) -> Result<(), String> { @@ -77,9 +66,10 @@ impl<'a> Cache<'a> { value: &T, ttl_secs: u64, ) -> Result<(), String> { - let mut conn = self.pool.get().await.map_err(|e| e.to_string())?; - let json = serde_json::to_string(value).map_err(|e| e.to_string())?; - conn.set_ex::<_, _, ()>(key, json, ttl_secs) + let json = serde_json::to_vec(value).map_err(|e| e.to_string())?; + let ttl = std::time::Duration::from_secs(ttl_secs); + let cache = self.cache().await.map_err(|e| e.to_string())?; + CacheTrait::set(cache, key, json, Some(ttl)) .await .map_err(|e| e.to_string())?; debug!("Cache: set key {} with TTL {}s", key, ttl_secs); @@ -87,18 +77,14 @@ impl<'a> Cache<'a> { } pub async fn delete(&self, key: &str) -> Result<(), String> { - let mut conn = self.pool.get().await.map_err(|e| e.to_string())?; - conn.del::<_, ()>(key).await.map_err(|e| e.to_string())?; - debug!("Cache: deleted key {}", key); - Ok(()) + let cache = self.cache().await.map_err(|e| e.to_string())?; + CacheTrait::invalidate(cache, key) + .await + .map_err(|e| e.to_string()) } pub async fn exists(&self, key: &str) -> bool { - let mut conn = match self.pool.get().await { - Ok(c) => c, - Err(_) => return false, - }; - conn.exists::<_, bool>(key).await.unwrap_or(false) + self.get::(key).await.is_some() } /// Get or set: returns cached value or computes and caches new value. diff --git a/src/infrastructure/repository/otakudesu.rs b/src/infrastructure/repository/otakudesu.rs index ebe9e90..e99af74 100644 --- a/src/infrastructure/repository/otakudesu.rs +++ b/src/infrastructure/repository/otakudesu.rs @@ -1,7 +1,6 @@ //! Otakudesu anime scraping repository. use async_trait::async_trait; -use backoff::future::retry; use tracing::{info, warn}; use crate::domain::entity::anime::{ @@ -13,7 +12,7 @@ use crate::domain::repository::ScrapingRepository; use crate::infrastructure::repository::parsers::otakudesu_parser; use crate::infrastructure::scraping::html_fetcher::fetch_html_with_retry; use crate::infrastructure::scraping::proxy_fetch::fetch_with_proxy; -use crate::infrastructure::scraping::retry::{default_backoff, transient}; +use crate::infrastructure::scraping::retry::{default_backoff, retry, retry_all}; const OTAKUDESU_BASE_URL: &str = "https://otakudesu.cloud"; @@ -208,14 +207,11 @@ impl OtakudesuRepository { } Err(e) => { warn!("Failed to fetch URL: {}, error: {:?}", url_owned, e); - Err(transient(ScrapingError::Http(format!( - "Proxy fetch failed: {}", - e - )))) + Err(ScrapingError::Http(format!("Proxy fetch failed: {}", e))) } } }; - retry(backoff, fetch_op) + retry(backoff, retry_all, fetch_op) .await .map_err(|e| ScrapingError::Http(e.to_string())) } diff --git a/src/infrastructure/scraping/parsing_utils.rs b/src/infrastructure/scraping/parsing_utils.rs index 685b5d1..906c535 100644 --- a/src/infrastructure/scraping/parsing_utils.rs +++ b/src/infrastructure/scraping/parsing_utils.rs @@ -4,8 +4,7 @@ use crate::domain::error::ScrapingError; use crate::infrastructure::scraping::proxy_fetch::fetch_with_proxy; -use crate::infrastructure::scraping::retry::{default_backoff, transient}; -use backoff::future::retry; +use crate::infrastructure::scraping::retry::{default_backoff, retry, retry_all}; use regex::Regex; use scraper::{ElementRef, Html, Selector}; use std::sync::LazyLock; @@ -23,12 +22,12 @@ pub async fn fetch_html_with_retry(url: &str) -> Result { } Err(e) => { warn!("Failed to fetch: {}, error: {:?}", url, e); - Err(transient(e)) + Err(e) } } }; - retry(backoff, fetch_operation) + retry(backoff, retry_all, fetch_operation) .await .map_err(|e| ScrapingError::Http(e.to_string())) } diff --git a/src/infrastructure/scraping/proxy_fetch.rs b/src/infrastructure/scraping/proxy_fetch.rs index a668e47..7c993aa 100644 --- a/src/infrastructure/scraping/proxy_fetch.rs +++ b/src/infrastructure/scraping/proxy_fetch.rs @@ -2,12 +2,11 @@ // Updated for sync Redis API, reqwest API changes, and concurrency optimization. use dashmap::DashMap; -use redis::AsyncCommands; use std::sync::LazyLock; use tokio::sync::broadcast; use tracing::{debug, error, warn}; -use crate::infrastructure::cache::redis_pool::get_redis_conn; +use crate::infrastructure::cache::mytheclipse; use crate::infrastructure::utils::cache_ttl::CACHE_TTL_VERY_SHORT; use crate::infrastructure::utils::http::common_headers; use crate::infrastructure::utils::http::is_internet_baik_block_page; @@ -52,13 +51,13 @@ fn get_fetch_cache_key(slug: &str) -> String { } async fn get_cached_fetch(slug: &str) -> Result, AppError> { - let mut conn = get_redis_conn().await?; let key = get_fetch_cache_key(slug); + let bytes = mytheclipse::get(&key) + .await + .map_err(|e| AppError::Internal(format!("Cache get failed for {}: {}", slug, e)))?; - let cached: Option = conn.get(&key).await?; - - if let Some(cached_str) = cached { - match serde_json::from_str::(&cached_str) { + if let Some(bytes) = bytes { + match serde_json::from_slice::(&bytes) { Ok(parsed) => { debug!("[fetchWithProxy] Returning cached response for {}", slug); Ok(Some(parsed)) @@ -71,14 +70,17 @@ async fn get_cached_fetch(slug: &str) -> Result, AppError> { } async fn set_cached_fetch(slug: &str, value: &FetchResult) -> Result<(), AppError> { - let mut conn = get_redis_conn().await?; let key = get_fetch_cache_key(slug); - let json_string = serde_json::to_string(value)?; + let json = serde_json::to_vec(value)?; // Use standardized TTL - conn.set_ex::<_, _, ()>(&key, &json_string, CACHE_TTL_VERY_SHORT) - .await?; - Ok(()) + mytheclipse::set( + &key, + json, + Some(std::time::Duration::from_secs(CACHE_TTL_VERY_SHORT)), + ) + .await + .map_err(|e| AppError::Internal(format!("Cache set failed for {}: {}", slug, e))) } // --- REDIS CACHE WRAPPER END --- diff --git a/src/infrastructure/scraping/retry.rs b/src/infrastructure/scraping/retry.rs index 6a40a1b..94c5880 100644 --- a/src/infrastructure/scraping/retry.rs +++ b/src/infrastructure/scraping/retry.rs @@ -1,16 +1,21 @@ -//! HTTP retry utilities with exponential backoff. +//! HTTP retry utilities — delegated to mytheclipse `retry` primitives. +//! +//! The scraper uses mytheclipse's retry machinery (`RetryConfig` + +//! `mytheclipse::retry`). These helpers keep the old backoff-style call +//! sites ergonomic while delegating the actual backoff/sleep/jitter logic +//! to the library. -use backoff::ExponentialBackoff; +use mytheclipse::{JitterKind, RetryConfig}; use std::time::Duration; /// Default retry configuration for HTTP requests. -pub fn default_backoff() -> ExponentialBackoff { - ExponentialBackoff { - initial_interval: Duration::from_millis(500), - max_interval: Duration::from_secs(10), - multiplier: 2.0, - max_elapsed_time: Some(Duration::from_secs(30)), - ..Default::default() +pub fn default_backoff() -> RetryConfig { + RetryConfig { + max_attempts: 4, + base_delay: Duration::from_millis(500), + max_delay: Duration::from_secs(10), + factor: 2.0, + jitter: JitterKind::Full, } } @@ -20,47 +25,102 @@ pub fn custom_backoff( max_secs: u64, multiplier: f64, max_elapsed_secs: u64, -) -> ExponentialBackoff { - ExponentialBackoff { - initial_interval: Duration::from_millis(initial_ms), - max_interval: Duration::from_secs(max_secs), - multiplier, - max_elapsed_time: Some(Duration::from_secs(max_elapsed_secs)), - ..Default::default() +) -> RetryConfig { + RetryConfig { + max_attempts: compute_attempts(initial_ms, max_secs, multiplier, max_elapsed_secs).max(1), + base_delay: Duration::from_millis(initial_ms), + max_delay: Duration::from_secs(max_secs), + factor: multiplier, + jitter: JitterKind::Full, } } +/// Estimate the number of attempts that fit in `max_elapsed_secs` given the +/// exponential backoff curve: solve `sum(base * factor^i) ≈ max_elapsed`. +fn compute_attempts(initial_ms: u64, max_secs: u64, multiplier: f64, max_elapsed_secs: u64) -> u32 { + if initial_ms == 0 || multiplier <= 1.0 { + return 1; + } + let mut elapsed_ms = 0u64; + let mut attempt = 0u64; + let max_ms = max_secs.saturating_mul(1000); + let budget_ms = max_elapsed_secs.saturating_mul(1000); + let mut delay_ms = initial_ms; + while elapsed_ms < budget_ms { + elapsed_ms = elapsed_ms.saturating_add(delay_ms); + attempt += 1; + delay_ms = ((delay_ms as f64) * multiplier).min(max_ms as f64) as u64; + } + attempt as u32 +} + /// Quick backoff for fast retries (3 attempts, 100ms initial). -pub fn quick_backoff() -> ExponentialBackoff { - ExponentialBackoff { - initial_interval: Duration::from_millis(100), - max_interval: Duration::from_secs(1), - multiplier: 2.0, - max_elapsed_time: Some(Duration::from_secs(5)), - ..Default::default() +pub fn quick_backoff() -> RetryConfig { + RetryConfig { + max_attempts: 3, + base_delay: Duration::from_millis(100), + max_delay: Duration::from_secs(1), + factor: 2.0, + jitter: JitterKind::Full, } } /// Slow backoff for long operations (10 attempts, 1s initial). -pub fn slow_backoff() -> ExponentialBackoff { - ExponentialBackoff { - initial_interval: Duration::from_secs(1), - max_interval: Duration::from_secs(30), - multiplier: 2.0, - max_elapsed_time: Some(Duration::from_secs(120)), - ..Default::default() +pub fn slow_backoff() -> RetryConfig { + RetryConfig { + max_attempts: 10, + base_delay: Duration::from_secs(1), + max_delay: Duration::from_secs(30), + factor: 2.0, + jitter: JitterKind::Full, } } -/// Make an error transient (will be retried). -pub fn transient(err: E) -> backoff::Error { - backoff::Error::transient(err) +/// Re-export mytheclipse retry for convenience. +pub use mytheclipse::retry; + +/// Re-export transient/permanent helpers for API compatibility. +/// +/// The legacy `backoff` crate distinguished transient vs permanent errors at +/// the error-type level; mytheclipse uses a retry predicate instead. All +/// scraped-HTTP failures are transient by nature (network/5xx), so both +/// helpers return the error unchanged and every call site retries everything +/// (`|_| true`). `permanent` is kept as a no-op alias for source +/// compatibility. +pub fn transient(e: E) -> E { + e } -/// Make an error permanent (will NOT be retried). -pub fn permanent(err: E) -> backoff::Error { - backoff::Error::permanent(err) +/// No-op alias for source compatibility (see [`transient`]). +pub fn permanent(e: E) -> E { + e } -// Re-export retry function for convenience -pub use backoff::future::retry; +/// Predicate used by all scraper retry loops: retry every error. +pub fn retry_all(_: &E) -> bool { + true +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn backoff_configs_build() { + assert_eq!(default_backoff().max_attempts, 4); + assert_eq!(quick_backoff().max_attempts, 3); + assert_eq!(slow_backoff().max_attempts, 10); + } + + #[test] + fn compute_attempts_curve() { + // 100ms base, 2x, 1s max, 5s budget → roughly 6 attempts. + let n = compute_attempts(100, 1, 2.0, 5); + assert!(n >= 4 && n <= 8, "got {n}"); + } + + #[test] + fn retry_all_retries() { + assert!(retry_all::(&std::io::Error::other("x"))); + } +} diff --git a/src/presentation/middleware/ratelimit.rs b/src/presentation/middleware/ratelimit.rs index 35e4e29..038f9ea 100644 --- a/src/presentation/middleware/ratelimit.rs +++ b/src/presentation/middleware/ratelimit.rs @@ -1,4 +1,8 @@ -//! Rate limiting middleware. +//! Rate limiting middleware — backed by `mytheclipse::RateLimiter`. +//! +//! The limiter is a token bucket from the mytheclipse core crate. The +//! middleware keeps the same axum shape (State>) and +//! `check()` semantics, but delegates token accounting to the library. use axum::{ extract::Request, @@ -7,42 +11,24 @@ use axum::{ response::{IntoResponse, Response}, Json, }; -use std::sync::atomic::{AtomicU64, Ordering}; -use std::sync::{Arc, Mutex}; -use std::time::Instant; +use mytheclipse::RateLimiter; +use std::sync::Arc; use crate::presentation::dto::common::ApiResponse; -/// Simple in-memory rate limiter. -pub struct RateLimiter { - max_requests: u64, - window_secs: u64, - counter: AtomicU64, - window_start: Mutex, +/// Convenience alias so callers don't need the mytheclipse import. +pub type AppRateLimiter = RateLimiter; + +/// Build a rate limiter (rate = requests/sec, burst = max burst capacity). +pub fn new_rate_limiter(rate_per_sec: f64, burst: u64) -> Arc { + Arc::new(RateLimiter::new(rate_per_sec, burst)) } -impl RateLimiter { - pub fn new(max_requests: u64, window_secs: u64) -> Arc { - Arc::new(Self { - max_requests, - window_secs, - counter: AtomicU64::new(0), - window_start: Mutex::new(Instant::now()), - }) - } - - pub fn check(&self) -> bool { - let Ok(mut window_guard) = self.window_start.lock() else { - return false; - }; - let window = &mut *window_guard; - if window.elapsed().as_secs() >= self.window_secs { - *window = Instant::now(); - self.counter.store(0, Ordering::SeqCst); - } - let count = self.counter.fetch_add(1, Ordering::SeqCst); - count < self.max_requests - } +/// Compatibility constructor matching the old (max_requests, window_secs) API. +/// Converts a fixed window into an equivalent token-bucket rate. +pub fn new_window_rate_limiter(max_requests: u64, window_secs: u64) -> Arc { + let rate_per_sec = max_requests as f64 / window_secs.max(1) as f64; + new_rate_limiter(rate_per_sec, max_requests) } pub async fn rate_limit_middleware( @@ -50,7 +36,7 @@ pub async fn rate_limit_middleware( request: Request, next: Next, ) -> Response { - if state.check() { + if state.try_acquire().is_ok() { next.run(request).await } else { (