From e24c139a50bf348a973e685dfc3f4e20632259ce Mon Sep 17 00:00:00 2001 From: Wayland Yang Date: Fri, 25 Sep 2026 17:49:59 +0800 Subject: [PATCH] Govern a base with one run at a time, so ten concurrent runs stop deciding the same cluster ten times Every document's extraction enqueues a govern job for its base, and enqueue_unless_queued only skips a job that is still queued; with 64 workers each new job starts at once. On the typed-graph bench nine govern runs overlapped on one base, all read the same queue heads, and 1346 review pairs received 8888 keep decisions (one pair ten decisions in 38 seconds from ten runs); a third of the run's model tokens went there, and the duplicate-key errors on agent_decisions_open were the concurrent proposals colliding. govern() now takes a per-base session-level advisory lock (try only, as the vector index build does). A run that does not get it exits and enqueues one govern for a minute later, deduplicated against queued jobs only, so pairs that arrive after the running job read its last queue head are still picked up. The queue and cluster reads skip rows the running job has marked adjudicating. close_review_auto reports how many rows it closed, and a keep or already-merged outcome records no decision when the row was closed by someone else. Store test: the second lock attempt on a base fails until the first is released; another base is another lock. Co-Authored-By: Claude Fable 5.1 Signed-off-by: Wayland Yang --- crates/utopia-server/src/governance.rs | 38 +++++++++++---- crates/utopia-store/src/governance.rs | 47 +++++++++++++++++-- crates/utopia-store/src/jobs.rs | 27 +++++++++++ crates/utopia-store/src/resolution.rs | 7 +-- .../tests/a_base_is_governed_by_one_run.rs | 35 ++++++++++++++ 5 files changed, 140 insertions(+), 14 deletions(-) create mode 100644 crates/utopia-store/tests/a_base_is_governed_by_one_run.rs 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(()) +}