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
38 changes: 30 additions & 8 deletions crates/utopia-server/src/governance.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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?;

Expand Down Expand Up @@ -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):合并会立刻送出图外的东西——违规、派生、答案——留给人,
Expand Down Expand Up @@ -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?;
Expand Down
47 changes: 44 additions & 3 deletions crates/utopia-store/src/governance.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<Postgres>,
key: String,
}

impl BaseLock {
pub async fn release(mut self) {
let unlocked: Result<bool, _> = 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<Option<BaseLock>> {
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<Vec<ReviewItem>> {
Expand Down
27 changes: 27 additions & 0 deletions crates/utopia-store/src/jobs.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Option<i64>> {
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<i64> {
enqueue_with_max_attempts(pool, kind, payload, 3).await
}
Expand Down
7 changes: 4 additions & 3 deletions crates/utopia-store/src/resolution.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<u64> {
let res = sqlx::query(
"UPDATE resolution_reviews SET status = $2, reason = $3, decided_at = now()
WHERE id = $1 AND status = 'pending'",
)
Expand All @@ -1680,7 +1681,7 @@ pub async fn close_review_auto(
.bind(reason)
.execute(pool)
.await?;
Ok(())
Ok(res.rows_affected())
}

/// 人工定夺。merge 方向:度数高(事实多)的一方作为存活目标,平局取更早创建的。
Expand Down
35 changes: 35 additions & 0 deletions crates/utopia-store/tests/a_base_is_governed_by_one_run.rs
Original file line number Diff line number Diff line change
@@ -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(())
}
Loading