diff --git a/crates/utopia-server/src/governance.rs b/crates/utopia-server/src/governance.rs index 34334da68..a2de57ce2 100644 --- a/crates/utopia-server/src/governance.rs +++ b/crates/utopia-server/src/governance.rs @@ -110,6 +110,20 @@ pub async fn govern(state: &AppState, kb_id: Uuid) -> anyhow::Result<()> { tracing::info!(%kb_id, "治理:没有配聊天模型,队列原地等"); return Ok(()); }; + // 一个库一次只跑一个治理任务。每篇文档抽完都排一个,而排队的去重只挡排着的、不挡在跑的: + // 抢不到锁就说明有人在治理这个库,它会把队列走完,走完还有积压会再排一个。它读完队头 + // 之后才进来的对它看不见,所以这里隔一分钟再排一个——排着的至多一个,跑着的也只有它 + let Some(base_lock) = gov::try_lock_base(&state.pool, kb_id).await? else { + tracing::info!(%kb_id, "治理:这个库已有任务在跑,一分钟后再看一眼"); + utopia_store::jobs::enqueue_unless_queued_after( + &state.pool, + "govern", + json!({ "kb_id": kb_id }), + std::time::Duration::from_secs(60), + ) + .await?; + return Ok(()); + }; let ctx = Ctx { state, kb_id, @@ -124,6 +138,7 @@ pub async fn govern(state: &AppState, kb_id: Uuid) -> anyhow::Result<()> { if let Err(e) = gov::release_locks(&state.pool, kb_id).await { tracing::warn!(%kb_id, error = %e, "治理:放锁失败"); } + base_lock.release().await; state.emit_review(kb_id); let more = outcome?; @@ -534,11 +549,14 @@ async fn apply(ctx: &Ctx<'_>, item: &ReviewItem, p: &Precedents, look: Look) -> utopia_store::resolution::survivor(pool, kb_id, item.right.id).await?, ); if l == r { - // 两边已经是同一个实体:只剩把审核行关上 - utopia_store::resolution::close_review_auto(pool, item.id, "merged", &reason) - .await?; - let id = gov::record(pool, kb_id, decision("applied", None)).await?; - audit(ctx, "review.merge", item, conf, id).await; + // 两边已经是同一个实体:只剩把审核行关上。关不上是已经有人关了,不再记一条 + let closed = + utopia_store::resolution::close_review_auto(pool, item.id, "merged", &reason) + .await?; + if closed > 0 { + let id = gov::record(pool, kb_id, decision("applied", None)).await?; + audit(ctx, "review.merge", item, conf, id).await; + } return Ok(()); } // 执行闸门(0027):合并会立刻送出图外的东西——违规、派生、答案——留给人, @@ -602,9 +620,13 @@ async fn apply(ctx: &Ctx<'_>, item: &ReviewItem, p: &Precedents, look: Look) -> } Gate::Apply => { let reason = format!("governed|{conf:.2}"); - utopia_store::resolution::close_review_auto(pool, item.id, "kept", &reason).await?; - let id = gov::record(pool, kb_id, decision("applied", None)).await?; - audit(ctx, "review.keep", item, conf, id).await; + // 关不上是这一对已经不是 pending(人裁了,或另一条路先到):不再记一条一样的裁决 + let closed = + utopia_store::resolution::close_review_auto(pool, item.id, "kept", &reason).await?; + if closed > 0 { + let id = gov::record(pool, kb_id, decision("applied", None)).await?; + audit(ctx, "review.keep", item, conf, id).await; + } } Gate::Propose => { utopia_store::resolution::escalate_review(pool, item.id, "proposed").await?; diff --git a/crates/utopia-store/src/governance.rs b/crates/utopia-store/src/governance.rs index 6476316c8..0eb8af84a 100644 --- a/crates/utopia-store/src/governance.rs +++ b/crates/utopia-store/src/governance.rs @@ -10,7 +10,7 @@ use chrono::{DateTime, Utc}; use serde::Serialize; -use sqlx::PgPool; +use sqlx::{PgPool, Postgres}; use utopia_core::models::{AgentDecisionView, ReviewItem}; use utopia_core::{AppError, AppResult}; use uuid::Uuid; @@ -615,9 +615,50 @@ pub fn settled_by_people( it.all(|x| x.merged() == first).then_some(first) } -/// 有一条开着的建议的对不进队列:agent 已经问过了,等人答 +/// 有一条开着的建议的对不进队列:agent 已经问过了,等人答。正在裁的(`adjudicating`) +/// 也不进:那是这一次任务自己手里的簇,读队头时它们还没落地 const OPEN_PROPOSAL: &str = "NOT EXISTS (SELECT 1 FROM agent_decisions d - WHERE d.target_kind = 'review' AND d.target_id = rr.id AND d.status = 'proposed')"; + WHERE d.target_kind = 'review' AND d.target_id = rr.id AND d.status = 'proposed') + AND rr.stage <> 'adjudicating'"; + +/// 一个库同一时刻只有一个治理任务在跑:会话级咨询锁,跟着这条连接走。 +/// +/// 每篇文档抽完都排一个治理任务,而 `jobs::enqueue_unless_queued` 只挡排着的、不挡在跑的: +/// 64 个 worker 把它们一起接起来,十个任务同时读同一个队头、同一簇裁十遍——一次 100 篇的 +/// 跑里 1346 个对被判了 8888 次,三分之一的 token 花在这上面。只用 **try**:抢不到就说明 +/// 有人在治理这个库,那个任务会把队列走完、有积压时再排一个 +const BASE_TRY_LOCK: &str = "SELECT pg_try_advisory_lock(hashtextextended('governance:' || $1, 0))"; +const BASE_UNLOCK: &str = "SELECT pg_advisory_unlock(hashtextextended('governance:' || $1, 0))"; + +/// 抢到的锁。放掉要显式调 [`BaseLock::release`];直接丢掉的话连接回池子时锁还挂着, +/// 所以 `release` 解不开就关连接,让 Postgres 收回它 +pub struct BaseLock { + conn: sqlx::pool::PoolConnection, + key: String, +} + +impl BaseLock { + pub async fn release(mut self) { + let unlocked: Result = sqlx::query_scalar(BASE_UNLOCK) + .bind(&self.key) + .fetch_one(&mut *self.conn) + .await; + if !matches!(unlocked, Ok(true)) { + let _ = self.conn.close().await; + } + } +} + +/// 试着拿这个库的治理锁。`None` = 别的任务正拿着 +pub async fn try_lock_base(pool: &PgPool, kb_id: Uuid) -> AppResult> { + let mut conn = pool.acquire().await?; + let key = kb_id.to_string(); + let got: bool = sqlx::query_scalar(BASE_TRY_LOCK) + .bind(&key) + .fetch_one(&mut *conn) + .await?; + Ok(got.then_some(BaseLock { conn, key })) +} /// 等人的重复对,先进先出 pub async fn queue(pool: &PgPool, kb_id: Uuid, limit: i64) -> AppResult> { diff --git a/crates/utopia-store/src/jobs.rs b/crates/utopia-store/src/jobs.rs index 1cb9db598..cfcc9f1c6 100644 --- a/crates/utopia-store/src/jobs.rs +++ b/crates/utopia-store/src/jobs.rs @@ -86,6 +86,33 @@ pub async fn enqueue_unless_queued( Ok(row.map(|(id,)| id)) } +/// 同 [`enqueue_unless_queued`],但晚一点跑。挡的只是排着的,不挡在跑的——调用方 +/// 正是那个在跑的任务、想给自己之后再排一个的时候,用这个而不是 `enqueue_unless_pending` +pub async fn enqueue_unless_queued_after( + pool: &PgPool, + kind: &str, + payload: serde_json::Value, + after: Duration, +) -> AppResult> { + let mut tx = pool.begin().await?; + let row: Option<(i64,)> = sqlx::query_as( + "INSERT INTO jobs (kind, payload, run_at) + SELECT $1, $2, now() + make_interval(secs => $3) + WHERE NOT EXISTS (SELECT 1 FROM jobs WHERE kind = $1 AND payload = $2 AND status = 'queued') + RETURNING id", + ) + .bind(kind) + .bind(payload) + .bind(after.as_secs_f64()) + .fetch_optional(&mut *tx) + .await?; + if row.is_some() { + notify_worker_tx(&mut tx).await?; + } + tx.commit().await?; + Ok(row.map(|(id,)| id)) +} + pub async fn enqueue(pool: &PgPool, kind: &str, payload: serde_json::Value) -> AppResult { enqueue_with_max_attempts(pool, kind, payload, 3).await } diff --git a/crates/utopia-store/src/resolution.rs b/crates/utopia-store/src/resolution.rs index fc4c706ae..49f25722c 100644 --- a/crates/utopia-store/src/resolution.rs +++ b/crates/utopia-store/src/resolution.rs @@ -1665,13 +1665,14 @@ pub async fn escalate_review(pool: &PgPool, review_id: Uuid, reason: &str) -> Ap } /// 自动定夺(LLM 高置信):merged / kept。合并动作本身由调用方先执行。 +/// 回关上了几行:0 = 这一对已经不是 pending(人裁了,或另一条路先到),调用方别再记一条一样的裁决 pub async fn close_review_auto( pool: &PgPool, review_id: Uuid, status: &str, reason: &str, -) -> AppResult<()> { - sqlx::query( +) -> AppResult { + let res = sqlx::query( "UPDATE resolution_reviews SET status = $2, reason = $3, decided_at = now() WHERE id = $1 AND status = 'pending'", ) @@ -1680,7 +1681,7 @@ pub async fn close_review_auto( .bind(reason) .execute(pool) .await?; - Ok(()) + Ok(res.rows_affected()) } /// 人工定夺。merge 方向:度数高(事实多)的一方作为存活目标,平局取更早创建的。 diff --git a/crates/utopia-store/tests/a_base_is_governed_by_one_run.rs b/crates/utopia-store/tests/a_base_is_governed_by_one_run.rs new file mode 100644 index 000000000..d2fd6f240 --- /dev/null +++ b/crates/utopia-store/tests/a_base_is_governed_by_one_run.rs @@ -0,0 +1,35 @@ +//! 一个库同一时刻只有一个治理任务(会话级咨询锁)。第二个任务抢不到锁就退出并晚点再排, +//! 而不是和第一个一起读同一个队头、把同一簇裁两遍。 + +use sqlx::PgPool; +use utopia_store::governance; +use uuid::Uuid; + +#[tokio::test] +async fn the_second_run_on_a_base_does_not_get_the_lock_until_the_first_releases_it( +) -> anyhow::Result<()> { + let Some(url) = utopia_store::test_db::url() else { + return Ok(()); + }; + let pool = PgPool::connect(&url).await?; + let kb = Uuid::now_v7(); + let other = Uuid::now_v7(); + + let first = governance::try_lock_base(&pool, kb) + .await? + .expect("a free base is locked at once"); + assert!( + governance::try_lock_base(&pool, kb).await?.is_none(), + "a second run on the same base must not get the lock" + ); + assert!( + governance::try_lock_base(&pool, other).await?.is_some(), + "another base is another lock" + ); + first.release().await; + let again = governance::try_lock_base(&pool, kb).await?; + assert!(again.is_some(), "released, the base can be governed again"); + again.unwrap().release().await; + pool.close().await; + Ok(()) +}