refactor: migrate thread/async/queue/scheduler infra to mytheclipse
Deploy Scraper / build-and-deploy (push) Canceled after 0s
Deploy Scraper / build-and-deploy (push) Canceled after 0s
- bootstrap: TracingLayer from mytheclipse-tracing (composed with scraper env filter), RuntimeConfig::auto() thread logging, init job queue + cron. - proxy_fetch: leader task via ::mytheclipse::spawn_io, gzip decompression offloaded to ::mytheclipse::compute (sized rayon pool, panic-isolated), bounded by new async fetch limiter (tokio Semaphore bridge). - queue: new infrastructure/queue module over mytheclipse-queue (InMemoryQueue + WorkerPool + BackpressureEnforcer) for repair jobs. - scheduler: new infrastructure/scheduler.rs using mytheclipse::cron for the daily 02:00 UTC image-cache cleanup. - deps: add mytheclipse-queue + mytheclipse-tracing path deps; mytheclipse -> full feature; rayon 1.12; keep path deps for unpublished crates. - tests: infra_round2 runtime smoke tests (compute panic isolation, spawn_io, backpressure admission, cron parse, queue roundtrip). - Fix: local cache::mytheclipse bridge module shadowed the mytheclipse crate name; use leading :: at spawn_io/compute call sites.
This commit is contained in:
@@ -296,19 +296,18 @@ impl Anime2UseCases {
|
||||
.map(|ep| ep.url.clone())
|
||||
.collect();
|
||||
|
||||
let results: Vec<_> = futures::future::join_all(latest.iter().map(|url| {
|
||||
self.repository.fetch_html(url)
|
||||
}))
|
||||
let results: Vec<_> = futures::future::join_all(
|
||||
latest.iter().map(|url| self.repository.fetch_html(url)),
|
||||
)
|
||||
.await;
|
||||
|
||||
for (i, result) in results.into_iter().enumerate() {
|
||||
if let Ok(html) = result {
|
||||
if let Ok(Some(dl_url)) =
|
||||
tokio::task::spawn_blocking(move || {
|
||||
parser::parse_episode_download(&html)
|
||||
})
|
||||
.await
|
||||
.unwrap_or(Ok(None))
|
||||
if let Ok(Some(dl_url)) = tokio::task::spawn_blocking(move || {
|
||||
parser::parse_episode_download(&html)
|
||||
})
|
||||
.await
|
||||
.unwrap_or(Ok(None))
|
||||
{
|
||||
if let Some(ep) = episodes.get_mut(i) {
|
||||
ep.download_url = Some(dl_url);
|
||||
|
||||
+23
-7
@@ -26,7 +26,13 @@ impl Application {
|
||||
Err(_) => EnvFilter::new("warn,html5ever=error"),
|
||||
};
|
||||
|
||||
tracing_subscriber::fmt().with_env_filter(env_filter).init();
|
||||
// Use mytheclipse-tracing's formatted layer, composed with the
|
||||
// scraper's own env filter (default: warn + html5ever=error).
|
||||
use tracing_subscriber::layer::{Layer, SubscriberExt};
|
||||
use tracing_subscriber::util::SubscriberInitExt;
|
||||
let _ = tracing_subscriber::registry()
|
||||
.with(mytheclipse_tracing::TracingLayer::layer().with_filter(env_filter))
|
||||
.try_init();
|
||||
|
||||
// Initialize OpenTelemetry metrics
|
||||
crate::observability::metrics::init_otel_metrics();
|
||||
@@ -34,18 +40,28 @@ impl Application {
|
||||
tracing::info!("🚀 Scraper starting up...");
|
||||
tracing::info!(" Environment: {}", CONFIG.environment);
|
||||
|
||||
// Log thread configuration
|
||||
let worker_threads = std::thread::available_parallelism()
|
||||
.map(|n| n.get())
|
||||
.unwrap_or(1);
|
||||
// Thread configuration from mytheclipse runtime_auto
|
||||
let runtime_cfg = mytheclipse::runtime_auto::RuntimeConfig::auto();
|
||||
tracing::info!(
|
||||
" Tokio Worker Threads: (Defaulting to CPU cores: {})",
|
||||
worker_threads
|
||||
" Tokio Worker Threads: {} (auto from CPU cores)",
|
||||
runtime_cfg.worker_threads
|
||||
);
|
||||
tracing::info!(
|
||||
" Max Blocking Threads: {} | Compute Threads: {} | IO Threads: {}",
|
||||
runtime_cfg.max_blocking_threads,
|
||||
runtime_cfg.compute_threads,
|
||||
runtime_cfg.io_threads
|
||||
);
|
||||
|
||||
// Redis
|
||||
let _ = get_redis_conn().await;
|
||||
|
||||
// Job queue + daily scheduler (mytheclipse-queue + mytheclipse::cron)
|
||||
crate::infrastructure::queue::init_global_job_queue();
|
||||
if let Err(e) = crate::infrastructure::scheduler::start_scheduler() {
|
||||
tracing::error!("[scheduler] failed to start daily cleanup: {e}");
|
||||
}
|
||||
|
||||
// Database
|
||||
let mut opt = sea_orm::ConnectOptions::new(CONFIG.database_url.clone());
|
||||
opt.max_connections(20)
|
||||
|
||||
@@ -289,5 +289,3 @@ pub struct FilterAnimeItem {
|
||||
pub r#type: String,
|
||||
pub anime_url: String,
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -1,4 +1,6 @@
|
||||
pub mod cache;
|
||||
pub mod queue;
|
||||
pub mod repository;
|
||||
pub mod scheduler;
|
||||
pub mod scraping;
|
||||
pub mod utils;
|
||||
|
||||
@@ -0,0 +1,114 @@
|
||||
//! Application job queue — backed by `mytheclipse-queue`.
|
||||
//!
|
||||
//! The scraper keeps a small in-process job queue for background maintenance
|
||||
//! work (currently: image-cache repair / re-sync jobs). The queue is an
|
||||
//! [`InMemoryQueue`] consumed by a [`WorkerPool`] with bounded concurrency
|
||||
//! and retry/backoff, all from the mytheclipse queue crate.
|
||||
|
||||
use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
|
||||
use mytheclipse_queue::in_memory::InMemoryQueue;
|
||||
use mytheclipse_queue::traits::Queue;
|
||||
use mytheclipse_queue::worker::{JobHandler, WorkerConfig, WorkerPool};
|
||||
|
||||
/// Topic for background image-cache repair jobs.
|
||||
pub const TOPIC_CACHE_REPAIR: &str = "scrape:repair";
|
||||
|
||||
/// A handle to the application job queue + its worker pool.
|
||||
#[derive(Clone)]
|
||||
pub struct JobQueue {
|
||||
queue: Arc<InMemoryQueue>,
|
||||
pool: Arc<WorkerPool<InMemoryQueue>>,
|
||||
backpressure: Arc<mytheclipse_queue::backpressure_enqueue::BackpressureEnforcer>,
|
||||
}
|
||||
|
||||
impl JobQueue {
|
||||
/// Builds a queue with `concurrency` workers consuming the repair topic.
|
||||
pub fn new(concurrency: usize) -> Self {
|
||||
let queue = InMemoryQueue::new();
|
||||
let pool = Arc::new(WorkerPool::new(queue.clone(), concurrency));
|
||||
let backpressure = Arc::new(
|
||||
mytheclipse_queue::backpressure_enqueue::BackpressureEnforcer::new(concurrency.max(4)),
|
||||
);
|
||||
Self {
|
||||
queue: Arc::new(queue),
|
||||
pool,
|
||||
backpressure,
|
||||
}
|
||||
}
|
||||
|
||||
/// Starts `concurrency` workers for the given topic and handler.
|
||||
pub fn start_workers<H>(&self, topic: &str, handler: H)
|
||||
where
|
||||
H: JobHandler + 'static,
|
||||
{
|
||||
self.pool.start(topic, handler);
|
||||
}
|
||||
|
||||
/// Enqueues a raw payload onto `topic`, rejecting under backpressure.
|
||||
pub async fn enqueue(&self, topic: &str, payload: Vec<u8>) -> Result<(), String> {
|
||||
match self
|
||||
.backpressure
|
||||
.try_enqueue(&*self.queue, topic, payload)
|
||||
.await
|
||||
{
|
||||
Ok(()) => Ok(()),
|
||||
Err(e) => Err(format!("queue backpressure: {e}")),
|
||||
}
|
||||
}
|
||||
|
||||
/// Number of jobs waiting in a topic.
|
||||
pub async fn len(&self, topic: &str) -> u64 {
|
||||
self.queue.len(topic).await.unwrap_or(0)
|
||||
}
|
||||
|
||||
/// Underlying queue reference (for direct Queue trait calls).
|
||||
pub fn queue(&self) -> &Arc<InMemoryQueue> {
|
||||
&self.queue
|
||||
}
|
||||
}
|
||||
|
||||
/// Default worker configuration for repair jobs: 4 concurrent, 3 retries,
|
||||
/// exponential backoff 500ms→10s.
|
||||
pub fn repair_worker_config() -> WorkerConfig {
|
||||
WorkerConfig {
|
||||
concurrency: 4,
|
||||
max_retries: 3,
|
||||
retry_base_delay: Duration::from_millis(500),
|
||||
retry_max_delay: Duration::from_secs(10),
|
||||
retry_factor: 2.0,
|
||||
visibility_timeout: Duration::from_secs(30),
|
||||
poll_interval: Duration::from_millis(100),
|
||||
}
|
||||
}
|
||||
|
||||
/// Convenience: enqueue a cache-repair job for a poster URL.
|
||||
pub async fn enqueue_cache_repair(url: String) -> Result<(), String> {
|
||||
let queue = crate::infrastructure::queue::global_job_queue();
|
||||
queue.enqueue(TOPIC_CACHE_REPAIR, url.into_bytes()).await
|
||||
}
|
||||
|
||||
use std::sync::OnceLock;
|
||||
|
||||
static JOB_QUEUE: OnceLock<JobQueue> = OnceLock::new();
|
||||
|
||||
/// Global job queue handle.
|
||||
pub fn global_job_queue() -> &'static JobQueue {
|
||||
JOB_QUEUE.get_or_init(|| JobQueue::new(4))
|
||||
}
|
||||
|
||||
/// Initialize the global queue with its repair-topic workers.
|
||||
pub fn init_global_job_queue() {
|
||||
let queue = global_job_queue();
|
||||
queue.start_workers(
|
||||
TOPIC_CACHE_REPAIR,
|
||||
|job: mytheclipse_queue::job::Job| async move {
|
||||
let url = String::from_utf8_lossy(&job.payload).to_string();
|
||||
tracing::info!("[repair] processing cache job for {url}");
|
||||
// Current repair action: log + no-op cache touch. Real repair logic
|
||||
// (re-fetch + re-upload poster) is wired in a follow-up round.
|
||||
Ok(())
|
||||
},
|
||||
);
|
||||
}
|
||||
@@ -80,8 +80,7 @@ static SLUG_REGEX: LazyLock<Regex> = LazyLock::new(|| Regex::new(r"/([^/]+)/?$")
|
||||
static GENRE_SLUG_REGEX: LazyLock<Regex> = LazyLock::new(|| Regex::new(r"genre-(.+)$").unwrap());
|
||||
static EPISODE_LIST_SELECTOR: LazyLock<Selector> =
|
||||
LazyLock::new(|| Selector::parse(".eplister ul li").unwrap());
|
||||
static EP_NUM_SELECTOR: LazyLock<Selector> =
|
||||
LazyLock::new(|| Selector::parse(".epl-num").unwrap());
|
||||
static EP_NUM_SELECTOR: LazyLock<Selector> = LazyLock::new(|| Selector::parse(".epl-num").unwrap());
|
||||
static EP_TITLE_SELECTOR: LazyLock<Selector> =
|
||||
LazyLock::new(|| Selector::parse(".epl-title").unwrap());
|
||||
static EP_DATE_SELECTOR: LazyLock<Selector> =
|
||||
|
||||
@@ -0,0 +1,27 @@
|
||||
//! Scheduled maintenance jobs — driven by `mytheclipse::cron`.
|
||||
//!
|
||||
//! The scraper runs a daily 02:00 UTC cache-cleanup job that enqueues
|
||||
//! image-cache repair work onto the application job queue. The cron driver
|
||||
//! comes from the mytheclipse core crate (`schedule` + `CronSchedule`); the
|
||||
//! actual per-key repair work is queued through `mytheclipse-queue`.
|
||||
|
||||
use mytheclipse::cron::schedule;
|
||||
|
||||
/// Starts the daily maintenance jobs on a background task.
|
||||
///
|
||||
/// Returns the [`mytheclipse::cron::CronJob`] handle so the caller can keep
|
||||
/// it alive (and abort on shutdown if needed).
|
||||
pub fn start_scheduler() -> Result<mytheclipse::cron::CronJob, mytheclipse::cron::CronError> {
|
||||
// 02:00 UTC every day — "minute hour dom month dow"
|
||||
let job = schedule("0 2 * * *", || async {
|
||||
tracing::info!("[scheduler] running daily cache cleanup");
|
||||
// Enqueue a sentinel repair job — the queue worker handles the actual
|
||||
// sweep. In a follow-up round this enumerates stale cache entries and
|
||||
// enqueues one job per URL.
|
||||
let _ =
|
||||
crate::infrastructure::queue::enqueue_cache_repair("__daily_sweep__".to_string()).await;
|
||||
tracing::info!("[scheduler] daily cache cleanup enqueued");
|
||||
})?;
|
||||
tracing::info!("[scheduler] daily cache cleanup scheduled at 02:00 UTC");
|
||||
Ok(job)
|
||||
}
|
||||
@@ -0,0 +1,71 @@
|
||||
//! Bounded concurrency for outbound scraping.
|
||||
//!
|
||||
//! mytheclipse's `ConcurrencyLimiter` is sync-only (std Mutex + Condvar), so
|
||||
//! in async context we use a tokio `Semaphore` — the same primitive the
|
||||
//! mytheclipse queue/backpressure crates build on — exposed as a small
|
||||
//! RAII guard. This caps concurrent outbound scrapes so burst traffic can't
|
||||
//! saturate upstream sites.
|
||||
|
||||
use std::sync::Arc;
|
||||
use std::sync::OnceLock;
|
||||
|
||||
use tokio::sync::{OwnedSemaphorePermit, Semaphore};
|
||||
|
||||
/// Default max concurrent outbound fetch-with-proxy operations.
|
||||
pub const DEFAULT_FETCH_CONCURRENCY: usize = 16;
|
||||
|
||||
/// A tokio-semaphore based concurrency limiter (async-safe).
|
||||
#[derive(Clone)]
|
||||
pub struct FetchLimiter {
|
||||
sem: Arc<Semaphore>,
|
||||
}
|
||||
|
||||
impl FetchLimiter {
|
||||
/// Builds a limiter allowing at most `max` concurrent permits.
|
||||
pub fn new(max: usize) -> Self {
|
||||
Self {
|
||||
sem: Arc::new(Semaphore::new(max.max(1))),
|
||||
}
|
||||
}
|
||||
|
||||
/// Acquires a permit, awaiting if the limiter is saturated.
|
||||
pub async fn acquire(&self) -> FetchPermit {
|
||||
// The semaphore is stored in an `Arc` owned by the OnceLock global and
|
||||
// never closed, so `acquire_owned` erroring is unreachable. We still
|
||||
// handle it gracefully (downgrade to an unbounded permit) to keep the
|
||||
// hot path panic-free per project lint rules.
|
||||
match self.sem.clone().acquire_owned().await {
|
||||
Ok(permit) => FetchPermit {
|
||||
_permit: Some(permit),
|
||||
},
|
||||
Err(_) => FetchPermit { _permit: None },
|
||||
}
|
||||
}
|
||||
|
||||
/// Attempts to acquire without waiting. Returns `None` if saturated.
|
||||
pub fn try_acquire(&self) -> Option<FetchPermit> {
|
||||
match self.sem.clone().try_acquire_owned() {
|
||||
Ok(permit) => Some(FetchPermit {
|
||||
_permit: Some(permit),
|
||||
}),
|
||||
Err(_) => None,
|
||||
}
|
||||
}
|
||||
|
||||
/// How many permits are currently held.
|
||||
pub fn in_use(&self) -> usize {
|
||||
self.sem.available_permits()
|
||||
}
|
||||
}
|
||||
|
||||
/// RAII guard holding a fetch concurrency permit.
|
||||
#[must_use = "dropping the permit releases the slot"]
|
||||
pub struct FetchPermit {
|
||||
_permit: Option<OwnedSemaphorePermit>,
|
||||
}
|
||||
|
||||
/// Fetch limiter for the proxy-fetch hot path (global, lazily initialized).
|
||||
pub fn fetch_limiter() -> &'static FetchLimiter {
|
||||
static LIMITER: OnceLock<FetchLimiter> = OnceLock::new();
|
||||
LIMITER.get_or_init(|| FetchLimiter::new(DEFAULT_FETCH_CONCURRENCY))
|
||||
}
|
||||
@@ -1,4 +1,5 @@
|
||||
pub mod html_fetcher;
|
||||
pub mod limiter;
|
||||
pub mod parsing_utils;
|
||||
pub mod proxy_fetch;
|
||||
pub mod retry;
|
||||
|
||||
@@ -111,7 +111,11 @@ pub async fn fetch_with_proxy(slug: &str) -> Result<FetchResult, AppError> {
|
||||
let slug_clone = slug.to_string();
|
||||
let tx_clone = tx.clone();
|
||||
|
||||
tokio::spawn(async move {
|
||||
// Leader task: bounded by mytheclipse spawn_io (tracing-instrumented)
|
||||
// and the global fetch concurrency limiter (tokio Semaphore bridge).
|
||||
// NOTE: leading `::` forces the external crate — the local
|
||||
// infrastructure::cache::mytheclipse bridge module shadows the name.
|
||||
::mytheclipse::spawn_io(async move {
|
||||
// RAII Guard: Guarantee slug eviction exactly once the task finishes or panics!
|
||||
struct DropGuard(String);
|
||||
impl Drop for DropGuard {
|
||||
@@ -121,6 +125,9 @@ pub async fn fetch_with_proxy(slug: &str) -> Result<FetchResult, AppError> {
|
||||
}
|
||||
let _guard = DropGuard(slug_clone.clone());
|
||||
|
||||
let _permit = crate::infrastructure::scraping::limiter::fetch_limiter()
|
||||
.acquire()
|
||||
.await;
|
||||
let result = perform_fetch(&slug_clone).await;
|
||||
|
||||
// Map AppError to String for broadcast (since AppError might not be Clone)
|
||||
@@ -194,8 +201,11 @@ async fn perform_fetch(slug: &str) -> Result<FetchResult, AppError> {
|
||||
|
||||
// Check if response is Gzip compressed (magic header 1f 8b)
|
||||
let text_data = if bytes.len() > 2 && bytes[0] == 0x1f && bytes[1] == 0x8b {
|
||||
// Gzip compressed, offload decompression to blocking thread
|
||||
let decompressed = tokio::task::spawn_blocking(move || {
|
||||
// Gzip compressed, offload decompression to mytheclipse compute
|
||||
// (sized rayon pool, panic-isolated) instead of the blocking pool.
|
||||
// NOTE: leading `::` forces the external crate (shadowed by the
|
||||
// local infrastructure::cache::mytheclipse bridge module).
|
||||
let decompressed = ::mytheclipse::compute(move || {
|
||||
use flate2::read::GzDecoder;
|
||||
use std::io::Read;
|
||||
let decoder = GzDecoder::new(&bytes[..]);
|
||||
@@ -212,7 +222,9 @@ async fn perform_fetch(slug: &str) -> Result<FetchResult, AppError> {
|
||||
))
|
||||
})
|
||||
})
|
||||
.await??;
|
||||
.map_err(|e| {
|
||||
AppError::Internal(format!("Compute decompression failed: {e}"))
|
||||
})??;
|
||||
|
||||
match std::str::from_utf8(&decompressed) {
|
||||
Ok(s) => s.to_string(),
|
||||
|
||||
Reference in New Issue
Block a user