From 4d347362940db00a71f22a9735013e401ccd6257 Mon Sep 17 00:00:00 2001 From: cybaka520 Date: Thu, 27 Aug 2026 23:17:00 +0800 Subject: [PATCH] =?UTF-8?q?refactor(worker):=20=E4=BC=98=E5=8C=96=E4=B8=8B?= =?UTF-8?q?=E8=BD=BD=E9=80=BB=E8=BE=91=E4=B8=8E=E5=90=8C=E6=AD=A5=E5=A4=B1?= =?UTF-8?q?=E8=B4=A5=E5=A4=84=E7=90=86?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- worker/src/sync/github.rs | 11 +++++++---- worker/src/worker/consumer.rs | 17 ++++++----------- worker/src/worker/sync_task.rs | 9 ++++++++- 3 files changed, 21 insertions(+), 16 deletions(-) diff --git a/worker/src/sync/github.rs b/worker/src/sync/github.rs index 6af680b..3c2416b 100644 --- a/worker/src/sync/github.rs +++ b/worker/src/sync/github.rs @@ -106,11 +106,14 @@ pub async fn download_raw_text(client: &Client, url: &str, cfg: &GitHubConfig) - .await } -/// 下载 TTML 文件原始字节,带 token +/// 下载 TTML 文件原始字节,带 token,失败无限重试直到成功 pub async fn download_raw_bytes(client: &Client, url: &str, cfg: &GitHubConfig) -> Result> { - let resp = send_github_get(client, url, None, cfg, "download bytes").await?; - let bytes = resp.bytes().await.context("read raw bytes")?; - Ok(bytes.to_vec()) + with_retry("download raw bytes", move || async move { + let resp = send_github_get(client, url, None, cfg, "download bytes").await?; + let bytes = resp.bytes().await.context("read raw bytes")?; + Ok(bytes.to_vec()) + }) + .await } /// 下载整包 zip diff --git a/worker/src/worker/consumer.rs b/worker/src/worker/consumer.rs index 93d28c6..14be30c 100644 --- a/worker/src/worker/consumer.rs +++ b/worker/src/worker/consumer.rs @@ -8,10 +8,10 @@ use crate::infra; use super::sync_task::SyncTaskRunner; -/// 同步任务最大重试次数 -const MAX_RETRIES: u32 = 3; /// 重试基础退避 const RETRY_BASE_DELAY: std::time::Duration = std::time::Duration::from_secs(5); +/// 重试退避上限(指数增长到此封顶) +const RETRY_MAX_DELAY: std::time::Duration = std::time::Duration::from_secs(60); /// 启动 RabbitMQ 消费循环 /// @@ -74,7 +74,8 @@ async fn handle_message( let runner = SyncTaskRunner::new(app.clone()); - // 应用层有限重试:处理 GitHub 限流、MinIO 抖动等瞬时故障 + // 应用层无限自动重试:处理 GitHub 限流、MinIO 抖动、下载失败等故障, + // 不成功不放弃,直到同步成功(退避指数增长,60s 封顶) let mut attempt: u32 = 0; loop { match runner.run(&request_id, &triggered_by, &payload).await { @@ -86,17 +87,11 @@ async fn handle_message( } Err(e) => { attempt += 1; - if attempt >= MAX_RETRIES { - // 最终失败由通用消费循环统一记录并 nack 进 DLQ, - // 这里附上 request_id 便于追踪 - return Err(e.context(format!("request_id={}, attempts={}", request_id, attempt))); - } - let delay = RETRY_BASE_DELAY * 2u32.saturating_pow(attempt - 1); + let delay = (RETRY_BASE_DELAY * 2u32.saturating_pow(attempt - 1)).min(RETRY_MAX_DELAY); warn!( request_id = %request_id, attempt, - max_retries = MAX_RETRIES, - delay_ms = delay.as_millis() as u64, + delay_secs = delay.as_secs(), error = %e, "sync attempt failed, will retry" ); diff --git a/worker/src/worker/sync_task.rs b/worker/src/worker/sync_task.rs index 609d2dd..ae7339c 100644 --- a/worker/src/worker/sync_task.rs +++ b/worker/src/worker/sync_task.rs @@ -250,9 +250,10 @@ impl SyncTaskRunner { info!("下载完成,成功下载 {} 个文件", downloaded.len()); // 5.5 强制 flush 最终下载进度到 DB(最终刷新直接 await,确保落库) + let final_failed; { let final_downloaded = progress_state.downloaded(); - let final_failed = progress_state.failed(); + final_failed = progress_state.failed(); debug!(final_downloaded, final_failed, "flush 最终下载进度"); run_progress_flush( repo_arc.clone(), @@ -264,6 +265,12 @@ impl SyncTaskRunner { .await; } + // 下载/上传存在失败文件时同步失败(不更新 last_synced_commit), + // 由消费者层重试;下次同步 diff 会重新包含这些文件,避免静默丢失 + if final_failed > 0 { + anyhow::bail!("{} 个文件下载/上传失败,本次同步标记失败以待重试", final_failed); + } + // 6. 并发解析 + 入库 + 累积 MeiliSearch 文档 let concurrency = self.app.cfg.worker.concurrency.max(1); let mut meili_docs: Vec = Vec::with_capacity(downloaded.len());