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
15 changes: 10 additions & 5 deletions crates/utopia-server/src/extraction_open.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1097,12 +1097,14 @@ pub(crate) async fn run_open(
tracing::info!(%document_id, "抽取任务已被新一轮接管,收尾时退出");
return Ok(());
}
// 有分块没抽成就不许标 done(与类型化那条路同一条规矩)
if let Some(msg) = incomplete_reason(&unextracted, chunks.len()) {
return Err(anyhow::anyhow!(msg));
// 有分块没抽成就不许标 done(与类型化那条路同一条规矩)。**但抽成了的那些照常收尾**:
// 从前在这里直接返回,后面的时间解析、待确认的那一声通知都跟着没了——记忆日志里一句
// 抽不成的话,让它后面每一句的确认卡都不出现(#1187)。错留到最后再报
let incomplete = incomplete_reason(&unextracted, chunks.len());
if incomplete.is_none() {
utopia_store::documents::set_graph_status(pool, document_id, "done").await?;
state.emit_document(kb_id, document_id);
}
utopia_store::documents::set_graph_status(pool, document_id, "done").await?;
state.emit_document(kb_id, document_id);
state.emit_graph(kb_id);
// 文档的时间语境(0064):这一轮各块报上来的日期并进文档上存着的那份。没抽到的块
// (增量抽取时认领的未变段落)原来报的留着;这一轮抽过的块以这一轮的为准。
Expand Down Expand Up @@ -1187,6 +1189,9 @@ pub(crate) async fn run_open(
if needs_adjudication || human_reviews_found {
state.emit_review(kb_id);
}
if let Some(msg) = incomplete {
return Err(anyhow::anyhow!(msg));
}
tracing::info!(%document_id, statements = statement_count, "开放图谱抽取完成");
Ok(())
}
Expand Down
92 changes: 92 additions & 0 deletions crates/utopia-server/src/extraction_open_out_of_credit_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -137,3 +137,95 @@ async fn an_empty_balance_stops_the_document_and_keeps_what_was_applied() -> any
assert_eq!(status, "failed");
Ok(())
}

/// #1187: one sentence of the memory log cannot be read. The attempt is still reported as
/// incomplete, but the sentence after it gets what an extracted sentence gets: its pending
/// statement and the time-resolution job. Before, the error returned ahead of both.
#[tokio::test]
async fn a_memory_that_cannot_be_extracted_does_not_hold_back_the_next() -> anyhow::Result<()> {
let Some(url) = utopia_store::test_db::url() else {
return Ok(());
};
let pool = sqlx::PgPool::connect(&url).await?;
utopia_store::db::migrate(&pool).await?;
let [org, ws, kb] = [(); 3].map(|_| Uuid::now_v7());
sqlx::raw_sql(&format!(
"INSERT INTO organizations(id,name) VALUES ('{org}','memory-incomplete');
INSERT INTO workspaces(id,org_id,name) VALUES ('{ws}','{org}','memory-incomplete');
INSERT INTO knowledge_bases(id,workspace_id,name) VALUES ('{kb}','{ws}','memory-incomplete');"
))
.execute(&pool)
.await?;
let now = chrono::Utc::now();
utopia_store::memory::append_episode(&pool, kb, "Something no model can read.", now).await?;
let (doc, second) =
utopia_store::memory::append_episode(&pool, kb, "Acme is based in London.", now).await?;
let good = json!({
"e": [["Acme", "organization", 1]],
"s": [["Acme is based in London.", "Acme", "based in", null, "London", null, null, null]],
"n": []
});
let script: Script = Arc::new(Mutex::new(vec![
Some("I cannot help with that.".to_string()),
Some(good.to_string()),
]));
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await?;
let endpoint = format!("http://{}", listener.local_addr()?);
let router = Router::new()
.route("/chat/completions", post(reply))
.with_state(script.clone());
tokio::spawn(async move { axum::serve(listener, router).await.unwrap() });
utopia_store::settings::upsert(
&pool,
ws,
Some(&endpoint),
None,
Some("scripted"),
None,
None,
None,
None,
)
.await?;
let dir = tempfile::tempdir()?;
let cfg = utopia_core::config::AppConfig {
data_dir: dir.path().to_string_lossy().into_owned(),
..Default::default()
};
let search = Arc::new(utopia_search::SearchIndex::open(
&dir.path().join("search"),
)?);
let state = AppState::new(pool.clone(), &cfg, search, "test-only".into());

let run = async {
let Err(err) = crate::extraction::extract_document(&state, doc, Proposer::default()).await
else {
anyhow::bail!("an attempt that left a sentence unread was reported as complete");
};
let (pending, resolving): (i64, i64) = sqlx::query_as(
"SELECT (SELECT count(*) FROM pending_facts WHERE chunk_id=$1),
(SELECT count(*) FROM jobs WHERE kind='resolve_time'
AND payload->>'document_id'=$2)",
)
.bind(second)
.bind(doc.to_string())
.fetch_one(&pool)
.await?;
Ok((format!("{err:#}"), pending, resolving))
}
.await;
sqlx::query("DELETE FROM jobs WHERE payload->>'document_id'=$1 OR payload->>'kb_id'=$2")
.bind(doc.to_string())
.bind(kb.to_string())
.execute(&pool)
.await?;
sqlx::query("DELETE FROM organizations WHERE id=$1")
.bind(org)
.execute(&pool)
.await?;
let (err, pending, resolving) = run?;
assert!(err.contains("1 of 2 chunks"), "{err}");
assert_eq!(pending, 1, "the second sentence awaits a nod");
assert_eq!(resolving, 1, "and its time words are queued for reading");
Ok(())
}
39 changes: 38 additions & 1 deletion crates/utopia-server/src/pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -273,8 +273,11 @@ pub async fn memory_ingest(
let kb_row = utopia_store::kbs::get(&state.pool, doc.kb_id).await?;
let settings = utopia_store::settings::get(&state.pool, kb_row.workspace_id).await?;

// 一句一句嵌,嵌不了的那句留着没有向量,别的照常往下走。日志是一篇越写越长的文档,
// 从前所有没向量的块一起嵌、一处出错整个任务失败:一句嵌不了的记忆从此挡在每一句
// 新记忆前面,索引不重建、抽取不排队,而 `remember` 照样回「记下了」(#1187)
if let Some((settings, client)) = embedder(settings.as_ref()) {
embed_pending(state, settings, &client, document_id).await?;
embed_each(state, settings, &client, document_id).await?;
}

let chunks = utopia_store::documents::chunks_full(&state.pool, document_id).await?;
Expand Down Expand Up @@ -304,6 +307,40 @@ pub async fn memory_ingest(
Ok(())
}

/// 记忆日志的嵌入:每块单独一次请求,失败的只记日志、留着下次再试,返回嵌成了几条。
/// 与 [`embed_pending`] 相反的取舍——那边一批错位会把向量写到别人身上,所以整篇失败;
/// 这里一次一条,没有错位可言,而一条失败不该连累后来的
async fn embed_each(
state: &AppState,
settings: &LlmSettings,
client: &LlmClient,
document_id: Uuid,
) -> anyhow::Result<usize> {
let pending =
utopia_store::documents::chunks_pending_embedding(&state.pool, document_id).await?;
let mut done = 0;
for (chunk_id, text) in pending {
let embedded = {
let _permit = llm_util::acquire_embed(state, settings).await;
client.embed(std::slice::from_ref(&text)).await
};
match embedded {
Ok(mut vectors) if vectors.len() == 1 => {
let vector = vectors.remove(0);
utopia_store::documents::set_embeddings(&state.pool, &[(chunk_id, vector)]).await?;
done += 1;
}
Ok(vectors) => {
tracing::warn!(%document_id, %chunk_id, got = vectors.len(), "记忆的向量数量对不上,这一句先不嵌");
}
Err(e) => {
tracing::warn!(%document_id, %chunk_id, error = %e, "这一句记忆没嵌成,别的照常");
}
}
}
Ok(done)
}

/// 工作区配了嵌入模型才有客户端;设置与客户端一起交出去,闸门许可证要按设置取。
fn embedder(settings: Option<&LlmSettings>) -> Option<(&LlmSettings, LlmClient)> {
let s = settings?;
Expand Down
33 changes: 31 additions & 2 deletions crates/utopia-server/src/pipeline_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@
//! 4. **一批失败文档不悬着**:走完整的 process_document,第二批回 500,文档落在 failed
//! 并带原因,不是停在 embedding。
//! 5. **就绪之前全部嵌完**:process_document 走通后没有一条向量为空。
//! 6. **记忆摄入走同一条路**:memory_ingest 嵌完自己的分块。
//! 6. **记忆一句一句嵌**(#1187):memory_ingest 嵌完自己的分块;一句嵌不了,别的照常。
//! 7. **正文夹 NUL 不毁整篇**(#611):Postgres 的 TEXT 不收 0x00,从前一个字节就让整篇
//! 落在 failed。现在走完整的 process_document 到 ready,库里没有一个分块带 NUL。
//!
Expand Down Expand Up @@ -1114,7 +1114,7 @@ async fn a_document_is_fully_embedded_before_it_is_ready() -> anyhow::Result<()>
}

#[tokio::test]
async fn a_memory_episode_embeds_by_the_same_path_as_a_document() -> anyhow::Result<()> {
async fn a_memory_episode_gets_the_vector_of_its_own_text() -> anyhow::Result<()> {
let Some(f) = fixture(FakeEmbed::new(Duration::ZERO)).await? else {
return Ok(());
};
Expand All @@ -1133,3 +1133,32 @@ async fn a_memory_episode_embeds_by_the_same_path_as_a_document() -> anyhow::Res
}
f.cleanup().await
}

/// #1187:从前日志里所有没向量的块一起嵌、一处出错整个任务失败,一句嵌不了的记忆从此
/// 挡在每一句新记忆前面。现在那一句留着没有向量,别的照常,任务照常走完
#[tokio::test]
async fn a_memory_that_cannot_be_embedded_does_not_hold_back_the_others() -> anyhow::Result<()> {
let mut fake = FakeEmbed::new(Duration::ZERO);
fake.fail_request = Some(2);
let Some(f) = fixture(fake).await? else {
return Ok(());
};
let doc = f.document_with_chunks(5).await?;
super::memory_ingest(
&f.state,
doc,
Proposer {
user_id: None,
token_id: None,
},
)
.await?;
let embedded: Vec<bool> = f
.stored(doc)
.await?
.iter()
.map(|(_, vector)| vector.is_some())
.collect();
assert_eq!(embedded, vec![true, false, true, true, true]);
f.cleanup().await
}
Loading