From 99f11804f5a10bcc82c56c68b0600a4383156e39 Mon Sep 17 00:00:00 2001 From: dada-yan Date: Mon, 21 Sep 2026 18:24:27 +1000 Subject: [PATCH 1/2] Report unconfirmed fallback answer persistence instead of success Signed-off-by: dada-yan --- .github/workflows/ci.yml | 6 + crates/utopia-server/src/api/chat.rs | 8 +- .../src/api/chat_empty_reply_tests.rs | 3 + .../src/api/chat_persistence_tests.rs | 113 ++++++++++++++++++ 4 files changed, 128 insertions(+), 2 deletions(-) create mode 100644 crates/utopia-server/src/api/chat_persistence_tests.rs diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 22fc4f962..9d237b900 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -100,6 +100,12 @@ jobs: UTOPIA_DATABASE_URL: postgres://utopia:utopia@localhost:5432/utopia UTOPIA_TEST_REQUIRE_DB: "1" + - name: Fallback answer persistence against Postgres + run: cargo test -p utopia-server chat_empty_reply_tests::persistence_tests + env: + UTOPIA_DATABASE_URL: postgres://utopia:utopia@localhost:5432/utopia + UTOPIA_TEST_REQUIRE_DB: "1" + - name: Hybrid retrieval against Postgres run: "cargo test -p utopia-server retrieval::" env: diff --git a/crates/utopia-server/src/api/chat.rs b/crates/utopia-server/src/api/chat.rs index 1b08111aa..785d9ec5d 100644 --- a/crates/utopia-server/src/api/chat.rs +++ b/crates/utopia-server/src/api/chat.rs @@ -944,7 +944,7 @@ fn legacy_rag( Err(e) => { yield error_event(&e.to_string()); return; } } } - let _ = utopia_store::conversations::append_message( + if let Err(error) = utopia_store::conversations::append_message( &state.pool, conversation_id, "assistant", &answer_acc, &utopia_store::conversations::TurnRecord { steps: serde_json::Value::Array(Vec::new()), @@ -952,7 +952,11 @@ fn legacy_rag( resolved: serde_json::Value::Array(Vec::new()), tool_exchange: serde_json::Value::Array(Vec::new()), }, - ).await; + ).await { + tracing::error!(%error, "fallback answer persistence was not confirmed"); + yield error_event("Could not confirm that the answer was saved."); + return; + } yield done_event(); } Err(e) => yield error_event(&e.to_string()), diff --git a/crates/utopia-server/src/api/chat_empty_reply_tests.rs b/crates/utopia-server/src/api/chat_empty_reply_tests.rs index 98d53fa8f..6ff144413 100644 --- a/crates/utopia-server/src/api/chat_empty_reply_tests.rs +++ b/crates/utopia-server/src/api/chat_empty_reply_tests.rs @@ -305,3 +305,6 @@ async fn a_reply_that_stays_empty_is_an_error_after_one_retry() -> anyhow::Resul ); f.cleanup().await } + +#[path = "chat_persistence_tests.rs"] +mod persistence_tests; diff --git a/crates/utopia-server/src/api/chat_persistence_tests.rs b/crates/utopia-server/src/api/chat_persistence_tests.rs new file mode 100644 index 000000000..ca42b7a67 --- /dev/null +++ b/crates/utopia-server/src/api/chat_persistence_tests.rs @@ -0,0 +1,113 @@ +//! The fallback answer can be streamed before its save; a failed save must still +//! terminate as error, including when the initiating browser has disconnected. +use super::*; + +#[derive(Clone, Default)] +struct LegacyOnly(Arc>>); +impl Respond for LegacyOnly { + fn respond(&self, request: &Request) -> ResponseTemplate { + let body: serde_json::Value = request.body_json().unwrap(); + let tools = body.get("tools").is_some(); + self.0.lock().unwrap().push(body); + if tools { + return ResponseTemplate::new(422) + .set_body_json(json!({"error":{"message":"tools unsupported"}})); + } + ResponseTemplate::new(200).insert_header("content-type","text/event-stream") + .set_body_string("data: {\"choices\":[{\"delta\":{\"content\":\"Generated answer.\"},\"finish_reason\":\"stop\"}]}\n\ndata: [DONE]\n\n") + } +} + +async fn exercise(deny: bool, disconnect: bool) -> anyhow::Result<()> { + let Some(f) = fixture(Scripted::new(vec![])).await? else { + return Ok(()); + }; + f._server.reset().await; + let model = LegacyOnly::default(); + Mock::given(method("POST")) + .and(path("/chat/completions")) + .respond_with(model.clone()) + .mount(&f._server) + .await; + let trigger = format!("reject_chat_{}", f.kb.simple()); + if deny { + sqlx::raw_sql(&format!("CREATE FUNCTION {trigger}() RETURNS trigger LANGUAGE plpgsql AS $$ BEGIN IF NEW.role='assistant' AND EXISTS(SELECT 1 FROM conversations WHERE id=NEW.conversation_id AND kb_id='{}') THEN RAISE EXCEPTION 'private persistence diagnostic'; END IF; RETURN NEW; END $$; CREATE TRIGGER {trigger} BEFORE INSERT ON conversation_messages FOR EACH ROW EXECUTE FUNCTION {trigger}();",f.kb)).execute(&f.pool).await?; + } + let result = tokio::time::timeout(std::time::Duration::from_secs(20), async { + let sse = if disconnect { + let id = + utopia_store::conversations::create(&f.pool, f.kb, f.user.id, "disconnect").await?; + let response = chat( + State(f.state.clone()), + AuthUser(f.user.clone()), + Path(f.kb), + Json(ChatReq { + conversation_id: Some(id), + message: "hello".into(), + }), + ) + .await + .map_err(|_| anyhow::anyhow!("chat refused"))?; + let (_, mut rx) = f + .state + .live + .attach(id) + .await + .ok_or_else(|| anyhow::anyhow!("producer ended before attachment"))?; + drop(response); + let mut frames = String::new(); + loop { + match rx.recv().await { + Ok(frame) => frames + .push_str(&format!("event: {}\ndata: {}\n\n", frame.event, frame.data)), + Err(tokio::sync::broadcast::error::RecvError::Closed) => break, + Err(e) => return Err(e.into()), + } + } + anyhow::ensure!(f.state.live.attach(id).await.is_none()); + frames + } else { + f.ask("hello").await? + }; + if deny { + anyhow::ensure!( + sse.contains("event: error") && !sse.contains("event: done"), + "{sse}" + ); + anyhow::ensure!(sse.contains("Could not confirm that the answer was saved.")); + anyhow::ensure!(!sse.contains("private persistence diagnostic")); + anyhow::ensure!(f.stored_answer().await?.is_none()); + } else { + anyhow::ensure!(sse.contains("event: done") && !sse.contains("event: error")); + anyhow::ensure!(f.stored_answer().await?.as_deref() == Some("Generated answer.")); + } + anyhow::ensure!( + model.0.lock().unwrap().len() == 3, + "compatibility negotiation plus exactly one answer, no save retry" + ); + Ok::<_, anyhow::Error>(()) + }) + .await; + if deny { + sqlx::raw_sql(&format!( + "DROP TRIGGER {trigger} ON conversation_messages; DROP FUNCTION {trigger}();" + )) + .execute(&f.pool) + .await?; + } + f.cleanup().await?; + result??; + Ok(()) +} +#[tokio::test] +async fn failed_fallback_save_never_reports_done() -> anyhow::Result<()> { + exercise(true, false).await +} +#[tokio::test] +async fn failed_fallback_save_after_disconnect_cleans_up() -> anyhow::Result<()> { + exercise(true, true).await +} +#[tokio::test] +async fn committed_fallback_answer_is_readable_at_done() -> anyhow::Result<()> { + exercise(false, false).await +} From 73f5442f50cc871415aa6be90df7d85b1ef71150 Mon Sep 17 00:00:00 2001 From: dada-yan Date: Mon, 21 Sep 2026 19:19:19 +1000 Subject: [PATCH 2/2] Use the shared PostgreSQL chat test step from PR 845 Signed-off-by: dada-yan --- .github/workflows/ci.yml | 6 ------ 1 file changed, 6 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 9d237b900..22fc4f962 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -100,12 +100,6 @@ jobs: UTOPIA_DATABASE_URL: postgres://utopia:utopia@localhost:5432/utopia UTOPIA_TEST_REQUIRE_DB: "1" - - name: Fallback answer persistence against Postgres - run: cargo test -p utopia-server chat_empty_reply_tests::persistence_tests - env: - UTOPIA_DATABASE_URL: postgres://utopia:utopia@localhost:5432/utopia - UTOPIA_TEST_REQUIRE_DB: "1" - - name: Hybrid retrieval against Postgres run: "cargo test -p utopia-server retrieval::" env: