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
98 changes: 98 additions & 0 deletions crates/utopia-server/src/api/jobs_routes.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,98 @@
//! 失败任务回队列(#216)。
//!
//! 一个失败任务原本只能通过它属的对象再跑:文档能重抽、来源能重同步;`bootstrap_ontology`、
//! `adjudicate_entities` 这些没有对象可点。余额耗尽(#201 让它第一次就 failed)一批文档
//! 全停,充值之后要逐个点。这里给两个入口:库内(Editor)与全局(管理员),范围可按
//! 种类与失败时间收窄——告警上的「再跑一遍」传的就是那次故障的时间窗。

use axum::extract::{Path, State};
use axum::Json;
use serde::Deserialize;
use serde_json::json;
use utopia_core::models::Role;
use utopia_store::jobs::RequeueScope;
use uuid::Uuid;

use super::graph_routes::require_kb;
use crate::auth::AuthUser;
use crate::error::ApiResult;
use crate::state::AppState;

#[derive(Deserialize, Default)]
pub struct RequeueBody {
#[serde(default)]
pub kind: Option<String>,
/// 只排这个时刻之后失败的
#[serde(default)]
pub failed_since: Option<chrono::DateTime<chrono::Utc>>,
}

pub async fn failed_in_kb(
State(state): State<AppState>,
AuthUser(user): AuthUser,
Path(kb_id): Path<Uuid>,
) -> ApiResult<Json<serde_json::Value>> {
require_kb(&state, &user, kb_id, Role::Viewer).await?;
let failed = utopia_store::jobs::failed_count(&state.pool, Some(kb_id)).await?;
Ok(Json(json!({ "failed": failed })))
}

pub async fn requeue_in_kb(
State(state): State<AppState>,
AuthUser(user): AuthUser,
Path(kb_id): Path<Uuid>,
Json(body): Json<RequeueBody>,
) -> ApiResult<Json<serde_json::Value>> {
require_kb(&state, &user, kb_id, Role::Editor).await?;
let requeued = utopia_store::jobs::requeue_failed(
&state.pool,
RequeueScope {
kb_id: Some(kb_id),
kind: body.kind.as_deref(),
failed_since: body.failed_since,
},
)
.await?;
let _ = utopia_store::audit::record(
&state.pool,
Some(kb_id),
user.id,
"jobs.requeued",
"kb",
Some(kb_id),
json!({ "requeued": requeued, "kind": body.kind, "failed_since": body.failed_since }),
)
.await;
Ok(Json(json!({ "requeued": requeued })))
}

/// 全局重排:系统级告警(没有库的)从这里走。只给管理员——它碰的是所有库的任务
pub async fn requeue_all(
State(state): State<AppState>,
AuthUser(user): AuthUser,
Json(body): Json<RequeueBody>,
) -> ApiResult<Json<serde_json::Value>> {
if !user.is_admin {
return Err(utopia_core::AppError::Forbidden.into());
}
let requeued = utopia_store::jobs::requeue_failed(
&state.pool,
RequeueScope {
kb_id: None,
kind: body.kind.as_deref(),
failed_since: body.failed_since,
},
)
.await?;
let _ = utopia_store::audit::record(
&state.pool,
None,
user.id,
"jobs.requeued",
"system",
None,
json!({ "requeued": requeued, "kind": body.kind, "failed_since": body.failed_since }),
)
.await;
Ok(Json(json!({ "requeued": requeued })))
}
5 changes: 5 additions & 0 deletions crates/utopia-server/src/api/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ mod datasource_routes;
mod documents_routes;
mod events_routes;
mod graph_routes;
mod jobs_routes;
mod kbs;
mod mapping_routes;
mod mcp;
Expand Down Expand Up @@ -86,6 +87,10 @@ pub fn router(state: AppState, cfg: &AppConfig) -> Router {
)
.route("/kbs/{id}/members", get(kbs::members))
.route("/kbs/{id}/audit", get(kbs::audit_log))
// 失败任务回队列(#216):库内给 Editor,全局给管理员
.route("/kbs/{id}/jobs/failed", get(jobs_routes::failed_in_kb))
.route("/kbs/{id}/jobs/requeue", post(jobs_routes::requeue_in_kb))
.route("/jobs/requeue", post(jobs_routes::requeue_all))
.route(
"/kbs/{id}/members/{user_id}",
axum::routing::put(kbs::set_member).delete(kbs::remove_member),
Expand Down
57 changes: 57 additions & 0 deletions crates/utopia-store/src/jobs.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,11 +4,13 @@
//! 并发消费:调度循环按"运行中 < 目标数"续派,任务在独立 task 执行;
//! 目标数经 AtomicUsize 热读——系统设置里改并发即时生效,无需重启。

use chrono::{DateTime, Utc};
use sqlx::PgPool;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
use std::time::Duration;
use utopia_core::AppResult;
use uuid::Uuid;

#[derive(Debug, Clone, sqlx::FromRow)]
pub struct Job {
Expand All @@ -29,6 +31,61 @@ pub async fn enqueue(pool: &PgPool, kind: &str, payload: serde_json::Value) -> A
Ok(id)
}

/// 重排失败任务的范围(#216)。三个条件都可空,空 = 不限。
///
/// **按库圈要解 payload**:任务表没有 kb 列,payload 只带 `document_id` /
/// `source_id` / `kb_id` 三种之一,各自解到库。没有库的系统任务只在不限库时才动
#[derive(Debug, Default, Clone, Copy)]
pub struct RequeueScope<'a> {
pub kb_id: Option<Uuid>,
pub kind: Option<&'a str>,
/// 只排这个时刻之后失败的——告警上的「再跑一遍」圈的正是那次故障窗口
pub failed_since: Option<DateTime<Utc>>,
}

/// 库范围的 SQL 谓词,`$N` 是库 id;`requeue_failed` 与 `failed_count` 共用
const KB_SCOPE: &str = "(
(j.payload ? 'kb_id' AND j.payload->>'kb_id' = $KB::text)
OR (j.payload ? 'document_id' AND EXISTS (
SELECT 1 FROM documents d
WHERE d.id::text = j.payload->>'document_id' AND d.kb_id = $KB))
OR (j.payload ? 'source_id' AND EXISTS (
SELECT 1 FROM sources s
WHERE s.id::text = j.payload->>'source_id' AND s.kb_id = $KB)))";

/// 把范围内的 failed 任务放回队列:`attempts` 归零、立即到期。
///
/// 处理器都是幂等的(启动时回收孤儿就靠这一点),所以重排永远安全;
/// 此前 `failed` 是终点,余额耗尽一批文档全失败,充值之后只能逐个点或整源重抽
pub async fn requeue_failed(pool: &PgPool, scope: RequeueScope<'_>) -> AppResult<u64> {
let sql = format!(
"UPDATE jobs j
SET status = 'queued', attempts = 0, run_at = now(), updated_at = now()
WHERE j.status = 'failed'
AND ($1::text IS NULL OR j.kind = $1)
AND ($2::timestamptz IS NULL OR j.updated_at >= $2)
AND ($3::uuid IS NULL OR {})",
KB_SCOPE.replace("$KB", "$3")
);
let res = sqlx::query(&sql)
.bind(scope.kind)
.bind(scope.failed_since)
.bind(scope.kb_id)
.execute(pool)
.await?;
Ok(res.rows_affected())
}

/// 范围内 failed 的条数——设置页那一行「N 个失败任务」
pub async fn failed_count(pool: &PgPool, kb_id: Option<Uuid>) -> AppResult<i64> {
let sql = format!(
"SELECT count(*) FROM jobs j
WHERE j.status = 'failed' AND ($1::uuid IS NULL OR {})",
KB_SCOPE.replace("$KB", "$1")
);
Ok(sqlx::query_scalar(&sql).bind(kb_id).fetch_one(pool).await?)
}

/// 认领一个到期任务;没有则返回 None。
async fn claim_one(pool: &PgPool) -> AppResult<Option<Job>> {
let job = sqlx::query_as(
Expand Down
Loading
Loading