Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 7 additions & 4 deletions worker/src/sync/github.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Vec<u8>> {
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
Expand Down
17 changes: 6 additions & 11 deletions worker/src/worker/consumer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 消费循环
///
Expand Down Expand Up @@ -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 {
Expand All @@ -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"
);
Expand Down
9 changes: 8 additions & 1 deletion worker/src/worker/sync_task.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(),
Expand All @@ -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<MeiliDocument> = Vec::with_capacity(downloaded.len());
Expand Down