mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-09 03:20:56 +00:00
fix(server): make generated title updates atomic
This commit is contained in:
parent
21e84484d2
commit
187e10879a
6 changed files with 83 additions and 21 deletions
|
|
@ -743,23 +743,6 @@ fn spawn_generated_title_task(task: GeneratedTitleTask) {
|
|||
return;
|
||||
}
|
||||
|
||||
let current = match task
|
||||
.state
|
||||
.stores
|
||||
.runs
|
||||
.get_cached_summary(&task.run_id, Utc::now())
|
||||
.await
|
||||
{
|
||||
Ok(Some(summary)) => summary,
|
||||
Ok(None) => return,
|
||||
Err(err) => {
|
||||
tracing::debug!(run_id = %task.run_id, error = %err, "Failed to re-read run summary for title update");
|
||||
return;
|
||||
}
|
||||
};
|
||||
if current.title != task.deterministic_title {
|
||||
return;
|
||||
}
|
||||
let run_store = match task.state.stores.runs.open_run(&task.run_id).await {
|
||||
Ok(store) => store,
|
||||
Err(err) => {
|
||||
|
|
@ -767,7 +750,8 @@ fn spawn_generated_title_task(task: GeneratedTitleTask) {
|
|||
return;
|
||||
}
|
||||
};
|
||||
if let Err(err) = workflow_event::append_event(
|
||||
let expected_title = task.deterministic_title;
|
||||
if let Err(err) = workflow_event::append_event_if(
|
||||
&run_store,
|
||||
&task.run_id,
|
||||
&workflow_event::Event::RunTitleUpdated {
|
||||
|
|
@ -776,6 +760,7 @@ fn spawn_generated_title_task(task: GeneratedTitleTask) {
|
|||
system_kind: SystemActorKind::Engine,
|
||||
}),
|
||||
},
|
||||
move |projection| projection.title().as_ref() == expected_title,
|
||||
)
|
||||
.await
|
||||
{
|
||||
|
|
|
|||
|
|
@ -3495,7 +3495,7 @@ async fn generated_title_does_not_overwrite_user_title_edit() {
|
|||
response_json!(response, StatusCode::OK).await;
|
||||
|
||||
wait_for_mock_hits(&title_mock, 1).await;
|
||||
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
|
||||
tokio::time::sleep(std::time::Duration::from_millis(250)).await;
|
||||
|
||||
assert_eq!(
|
||||
state
|
||||
|
|
|
|||
|
|
@ -803,6 +803,40 @@ mod tests {
|
|||
assert!(matches!(err, Error::ReadOnly));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn append_event_if_evaluates_latest_projection_before_appending() {
|
||||
let (_object_store, store) = make_store();
|
||||
let run = store.create_run(&test_run_id("run-1")).await.unwrap();
|
||||
append_created(&run, "run-1", dt("2026-03-27T12:00:00Z")).await;
|
||||
let initial_title = run.state().await.unwrap().title().into_owned();
|
||||
|
||||
run.append_event(&event_payload(
|
||||
"run-1",
|
||||
"2026-03-27T12:00:01Z",
|
||||
"run.title.updated",
|
||||
&serde_json::json!({ "title": "User title" }),
|
||||
))
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let generated_update = event_payload(
|
||||
"run-1",
|
||||
"2026-03-27T12:00:02Z",
|
||||
"run.title.updated",
|
||||
&serde_json::json!({ "title": "Generated title" }),
|
||||
);
|
||||
let appended = run
|
||||
.append_event_if(&generated_update, |projection| {
|
||||
projection.title() == initial_title
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(appended, None);
|
||||
assert_eq!(run.state().await.unwrap().title(), "User title");
|
||||
assert_eq!(run.list_events().await.unwrap().len(), 2);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn control_request_events_set_pending_control_without_overwriting_status() {
|
||||
let (_object_store, store) = make_store();
|
||||
|
|
|
|||
|
|
@ -169,6 +169,10 @@ impl RunDatabase {
|
|||
|
||||
async fn projected_state(&self) -> Result<RunProjection> {
|
||||
let _state_guard = self.inner.state_lock.lock().await;
|
||||
self.projected_state_locked().await
|
||||
}
|
||||
|
||||
async fn projected_state_locked(&self) -> Result<RunProjection> {
|
||||
let next_seq = {
|
||||
let cache = self.inner.projection_cache.lock().await;
|
||||
cache.last_seq.saturating_add(1)
|
||||
|
|
@ -254,12 +258,35 @@ impl RunDatabase {
|
|||
Ok(self.append_event_envelope(payload).await?.seq)
|
||||
}
|
||||
|
||||
/// Atomically appends `payload` when `predicate` matches the latest run
|
||||
/// projection.
|
||||
pub async fn append_event_if(
|
||||
&self,
|
||||
payload: &EventPayload,
|
||||
predicate: impl FnOnce(&RunProjection) -> bool,
|
||||
) -> Result<Option<u32>> {
|
||||
if self.read_only {
|
||||
return Err(Error::ReadOnly);
|
||||
}
|
||||
payload.validate(&self.inner.run_id)?;
|
||||
let _state_guard = self.inner.state_lock.lock().await;
|
||||
let projection = self.projected_state_locked().await?;
|
||||
if !predicate(&projection) {
|
||||
return Ok(None);
|
||||
}
|
||||
Ok(Some(self.append_event_envelope_locked(payload).await?.seq))
|
||||
}
|
||||
|
||||
pub async fn append_event_envelope(&self, payload: &EventPayload) -> Result<EventEnvelope> {
|
||||
if self.read_only {
|
||||
return Err(Error::ReadOnly);
|
||||
}
|
||||
payload.validate(&self.inner.run_id)?;
|
||||
let _state_guard = self.inner.state_lock.lock().await;
|
||||
self.append_event_envelope_locked(payload).await
|
||||
}
|
||||
|
||||
async fn append_event_envelope_locked(&self, payload: &EventPayload) -> Result<EventEnvelope> {
|
||||
let seq = self.inner.event_seq.fetch_add(1, Ordering::SeqCst);
|
||||
let event = EventEnvelope {
|
||||
seq,
|
||||
|
|
|
|||
|
|
@ -18,6 +18,7 @@ pub use self::redaction::{
|
|||
build_redacted_event_payload, event_payload_from_redacted_json, redacted_event_json,
|
||||
};
|
||||
pub use self::sink::{
|
||||
RunEventLogger, RunEventSink, StoreProgressLogger, append_event, append_event_to_sink,
|
||||
RunEventLogger, RunEventSink, StoreProgressLogger, append_event, append_event_if,
|
||||
append_event_to_sink,
|
||||
};
|
||||
pub use crate::stage_scope::StageScope;
|
||||
|
|
|
|||
|
|
@ -2,7 +2,7 @@ use std::future::Future;
|
|||
use std::pin::Pin;
|
||||
use std::sync::Arc;
|
||||
|
||||
use ::fabro_types::{RunEvent, RunId};
|
||||
use ::fabro_types::{RunEvent, RunId, RunProjection};
|
||||
use anyhow::Result;
|
||||
use fabro_store::RunDatabase;
|
||||
use tokio::io::{AsyncWrite, AsyncWriteExt};
|
||||
|
|
@ -23,6 +23,21 @@ pub async fn append_event(run_store: &RunDatabase, run_id: &RunId, event: &Event
|
|||
.map_err(anyhow::Error::from)
|
||||
}
|
||||
|
||||
pub async fn append_event_if(
|
||||
run_store: &RunDatabase,
|
||||
run_id: &RunId,
|
||||
event: &Event,
|
||||
predicate: impl FnOnce(&RunProjection) -> bool,
|
||||
) -> Result<bool> {
|
||||
let stored = to_run_event(run_id, event);
|
||||
let payload = build_redacted_event_payload(&stored, run_id)?;
|
||||
run_store
|
||||
.append_event_if(&payload, predicate)
|
||||
.await
|
||||
.map(|seq| seq.is_some())
|
||||
.map_err(anyhow::Error::from)
|
||||
}
|
||||
|
||||
pub async fn append_event_to_sink(
|
||||
sink: &RunEventSink,
|
||||
run_id: &RunId,
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue