feat: collect flow conversations and agent memory once their last message goes (#11178)

* feat: collect flow conversations and agent memory when their last message goes

* fix: lock the conversation lookup so a new turn orders against its cleanup

* fix: let concurrent turns recreate a collected conversation without conflicting
This commit is contained in:
Guilhem
2026-09-17 09:27:12 +02:00
committed by GitHub
parent a9ec0aec3a
commit 23c24a9688
18 changed files with 614 additions and 59 deletions
@@ -0,0 +1,15 @@
{
"db_name": "PostgreSQL",
"query": "DELETE FROM flow_conversation c\n WHERE c.id = ANY($1)\n AND c.workspace_id = $2\n AND NOT EXISTS (\n SELECT 1 FROM flow_conversation_message m WHERE m.conversation_id = c.id\n )",
"describe": {
"columns": [],
"parameters": {
"Left": [
"UuidArray",
"Text"
]
},
"nullable": []
},
"hash": "07a005f0f9e80a156cd2a5a0ae39a1fabeaa167818206a25abfe31d5582f942a"
}
@@ -1,15 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "DELETE FROM flow_conversation_message m\n USING flow_conversation c\n WHERE m.conversation_id = c.id AND c.workspace_id = $1 AND m.job_id = ANY($2)",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Text",
"UuidArray"
]
},
"nullable": []
},
"hash": "462d2b2822b185a6f51fafcfa957cb3b31ee6b69abae79a214dddba0dee4425c"
}
@@ -0,0 +1,23 @@
{
"db_name": "PostgreSQL",
"query": "SELECT id FROM flow_conversation WHERE id = ANY($1) AND workspace_id = $2 ORDER BY id FOR UPDATE",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "id",
"type_info": "Uuid"
}
],
"parameters": {
"Left": [
"UuidArray",
"Text"
]
},
"nullable": [
false
]
},
"hash": "4f52bf546579f26a1d22c239b8b0054b753cfbeb5dad1e8120fd5e8a672d50ef"
}
@@ -0,0 +1,23 @@
{
"db_name": "PostgreSQL",
"query": "DELETE FROM flow_conversation_message m\n USING flow_conversation c\n WHERE m.conversation_id = c.id AND c.workspace_id = $1 AND m.job_id = ANY($2)\n RETURNING m.conversation_id",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "conversation_id",
"type_info": "Uuid"
}
],
"parameters": {
"Left": [
"Text",
"UuidArray"
]
},
"nullable": [
false
]
},
"hash": "69bfbe9b39414b724488532cc3b3659915d9fcb3d58f16532aaffe06c44ec976"
}
@@ -0,0 +1,59 @@
{
"db_name": "PostgreSQL",
"query": "SELECT id, workspace_id, flow_path, title, created_at, updated_at, created_by\n FROM flow_conversation\n WHERE id = $1 AND workspace_id = $2\n FOR UPDATE",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "id",
"type_info": "Uuid"
},
{
"ordinal": 1,
"name": "workspace_id",
"type_info": "Varchar"
},
{
"ordinal": 2,
"name": "flow_path",
"type_info": "Varchar"
},
{
"ordinal": 3,
"name": "title",
"type_info": "Varchar"
},
{
"ordinal": 4,
"name": "created_at",
"type_info": "Timestamptz"
},
{
"ordinal": 5,
"name": "updated_at",
"type_info": "Timestamptz"
},
{
"ordinal": 6,
"name": "created_by",
"type_info": "Varchar"
}
],
"parameters": {
"Left": [
"Uuid",
"Text"
]
},
"nullable": [
false,
false,
false,
true,
false,
false,
false
]
},
"hash": "6f32c1feed096ff706ae359ad6a3ca33b3f82ca38289dfa4a69aa95041027d57"
}
@@ -0,0 +1,15 @@
{
"db_name": "PostgreSQL",
"query": "DELETE FROM ai_agent_memory a\n USING flow_conversation c\n WHERE c.id = ANY($1)\n AND c.workspace_id = $2\n AND a.conversation_id = c.id\n AND a.workspace_id = c.workspace_id\n AND NOT EXISTS (\n SELECT 1 FROM flow_conversation_message m WHERE m.conversation_id = c.id\n )",
"describe": {
"columns": [],
"parameters": {
"Left": [
"UuidArray",
"Text"
]
},
"nullable": []
},
"hash": "89ea81b765550cf665e30533efc9672f8c98d2579fc72c07752361cc5fd683dc"
}
@@ -0,0 +1,14 @@
{
"db_name": "PostgreSQL",
"query": "DELETE FROM flow_conversation c\n WHERE c.id = ANY($1)\n AND NOT EXISTS (\n SELECT 1 FROM flow_conversation_message m WHERE m.conversation_id = c.id\n )",
"describe": {
"columns": [],
"parameters": {
"Left": [
"UuidArray"
]
},
"nullable": []
},
"hash": "90b910da8d00a7c7bcf29c167e38e44eb1c0062a8241dc8fe0ed3dd95b65f89a"
}
@@ -0,0 +1,22 @@
{
"db_name": "PostgreSQL",
"query": "DELETE FROM flow_conversation_message WHERE job_id = ANY($1) RETURNING conversation_id",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "conversation_id",
"type_info": "Uuid"
}
],
"parameters": {
"Left": [
"UuidArray"
]
},
"nullable": [
false
]
},
"hash": "967f52005f4a044b3a2e9f02ceadf90dad5246681bde6caaa633b85a5e8b2352"
}
@@ -0,0 +1,14 @@
{
"db_name": "PostgreSQL",
"query": "DELETE FROM ai_agent_memory a\n USING flow_conversation c\n WHERE c.id = ANY($1)\n AND a.conversation_id = c.id\n AND a.workspace_id = c.workspace_id\n AND NOT EXISTS (\n SELECT 1 FROM flow_conversation_message m WHERE m.conversation_id = c.id\n )",
"describe": {
"columns": [],
"parameters": {
"Left": [
"UuidArray"
]
},
"nullable": []
},
"hash": "a1b23f3e62c6433d95cdac58215741a51ca0ce66bf1c674097705cf2e1b72eff"
}
@@ -1,14 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "DELETE FROM flow_conversation_message WHERE job_id = ANY($1)",
"describe": {
"columns": [],
"parameters": {
"Left": [
"UuidArray"
]
},
"nullable": []
},
"hash": "bfdd60b42e32bd81e2d20b327462893147b4e5ff078531de36147d908132d636"
}
@@ -1,6 +1,6 @@
{
"db_name": "PostgreSQL",
"query": "INSERT INTO flow_conversation (id, workspace_id, flow_path, created_by, title)\n VALUES ($1, $2, $3, $4, $5)\n RETURNING id, workspace_id, flow_path, title, created_at, updated_at, created_by",
"query": "INSERT INTO flow_conversation (id, workspace_id, flow_path, created_by, title)\n VALUES ($1, $2, $3, $4, $5)\n ON CONFLICT (id) DO NOTHING\n RETURNING id, workspace_id, flow_path, title, created_at, updated_at, created_by",
"describe": {
"columns": [
{
@@ -58,5 +58,5 @@
false
]
},
"hash": "6bd23a98838e3eec309e6b696edc776bd56fc9dae1238b3272557d1562400dbe"
"hash": "c1e3ed3ecc3bcb98f60ba8196d33fee4a74f61b061e5025ecb75882208b3ba8f"
}
@@ -0,0 +1,22 @@
{
"db_name": "PostgreSQL",
"query": "SELECT id FROM flow_conversation WHERE id = ANY($1) ORDER BY id FOR UPDATE",
"describe": {
"columns": [
{
"ordinal": 0,
"name": "id",
"type_info": "Uuid"
}
],
"parameters": {
"Left": [
"UuidArray"
]
},
"nullable": [
false
]
},
"hash": "ec295b3890a0018475ec0a3774c7a30d71a5689efe72daf58bd1e8f6cf90c410"
}
+247 -1
View File
@@ -37,7 +37,7 @@ async fn seed_side_rows(db: &Pool<Postgres>, ws: &str, job_id: Uuid) -> anyhow::
.bind(ws)
.execute(db)
.await?;
// created_seq is assigned by a trigger; inserting a value is rejected.
// created_seq is an identity column; supplying a value is rejected.
sqlx::query(
"INSERT INTO flow_conversation_message (conversation_id, message_type, content, job_id)
VALUES ($1, 'assistant', 'hi', $2)",
@@ -121,6 +121,184 @@ async fn test_delete_jobs_removes_side_rows(db: Pool<Postgres>) -> anyhow::Resul
Ok(())
}
/// (conversation rows, agent-memory rows) for one conversation.
async fn conversation_and_memory_counts(
db: &Pool<Postgres>,
conversation_id: Uuid,
) -> anyhow::Result<(i64, i64)> {
Ok((
count(
db,
"SELECT count(*) FROM flow_conversation WHERE id = $1",
conversation_id,
)
.await?,
count(
db,
"SELECT count(*) FROM ai_agent_memory WHERE conversation_id = $1",
conversation_id,
)
.await?,
))
}
/// A conversation outlives the jobs behind its messages until the last one goes: only then
/// are the row and the agent's memory for it left with nothing, and only then are they
/// deleted. Both halves matter — the surviving half is what a single data-modifying CTE
/// would break, since its emptiness check would read the snapshot from before the delete.
#[sqlx::test(fixtures("base"))]
async fn test_delete_jobs_removes_a_conversation_once_its_last_message_goes(
db: Pool<Postgres>,
) -> anyhow::Result<()> {
initialize_tracing().await;
let first_job = Uuid::new_v4();
let second_job = Uuid::new_v4();
insert_job(&db, WS, first_job).await?;
insert_job(&db, WS, second_job).await?;
let conv_id = Uuid::new_v4();
sqlx::query(
"INSERT INTO flow_conversation (id, workspace_id, flow_path, created_by)
VALUES ($1, $2, 'f/flow', 'test-user')",
)
.bind(conv_id)
.bind(WS)
.execute(&db)
.await?;
for job_id in [first_job, second_job] {
sqlx::query(
"INSERT INTO flow_conversation_message (conversation_id, message_type, content, job_id)
VALUES ($1, 'assistant', 'hi', $2)",
)
.bind(conv_id)
.bind(job_id)
.execute(&db)
.await?;
}
sqlx::query(
"INSERT INTO ai_agent_memory (workspace_id, conversation_id, step_id, messages)
VALUES ($1, $2, 'a', '[]'::jsonb)",
)
.bind(WS)
.bind(conv_id)
.execute(&db)
.await?;
let mut conn = db.acquire().await?;
windmill_common::jobs::delete_jobs(&mut conn, &[first_job]).await?;
drop(conn);
assert_eq!(
conversation_and_memory_counts(&db, conv_id).await?,
(1, 1),
"a conversation with a message left must survive, memory included"
);
let mut conn = db.acquire().await?;
windmill_common::jobs::delete_jobs(&mut conn, &[second_job]).await?;
drop(conn);
assert_eq!(
conversation_and_memory_counts(&db, conv_id).await?,
(0, 0),
"the last message going should take the conversation and its memory"
);
Ok(())
}
/// Turns that start while retention is collecting their conversation must land, not fail:
/// the conversation lookup locks the row, so each turn waits for the collector's commit,
/// finds the conversation gone, and creates it again — the first insert wins and the other
/// reads its row. Without the lock a turn's message insert is what waits, on the parent
/// row's key lock, and fails its FK check afterwards.
#[sqlx::test(fixtures("base"))]
async fn test_new_turns_wait_for_conversation_cleanup_and_recreate(
db: Pool<Postgres>,
) -> anyhow::Result<()> {
initialize_tracing().await;
let old_job = Uuid::new_v4();
let new_jobs = [Uuid::new_v4(), Uuid::new_v4()];
insert_job(&db, WS, old_job).await?;
for job in new_jobs {
insert_job(&db, WS, job).await?;
}
let conv_id = Uuid::new_v4();
sqlx::query(
"INSERT INTO flow_conversation (id, workspace_id, flow_path, created_by)
VALUES ($1, $2, 'f/flow', 'test-user')",
)
.bind(conv_id)
.bind(WS)
.execute(&db)
.await?;
sqlx::query(
"INSERT INTO flow_conversation_message (conversation_id, message_type, content, job_id)
VALUES ($1, 'user', 'hi', $2)",
)
.bind(conv_id)
.bind(old_job)
.execute(&db)
.await?;
// The collector holds the conversation row locked and deleted, uncommitted.
let mut cleanup = db.begin().await?;
windmill_common::jobs::delete_jobs(&mut *cleanup, &[old_job]).await?;
let turns: Vec<_> = new_jobs
.into_iter()
.map(|new_job| {
let db = db.clone();
tokio::spawn(async move {
let mut tx = db.begin().await?;
windmill_common::flow_conversations::get_or_create_conversation_with_id(
&mut tx,
WS,
"f/flow",
"test-user",
"hi again",
conv_id,
)
.await?;
windmill_common::flow_conversations::add_message_to_conversation_tx(
&mut tx,
conv_id,
Some(new_job),
"hi again",
windmill_common::flow_conversations::MessageType::User,
None,
true,
)
.await?;
tx.commit().await?;
anyhow::Ok(())
})
})
.collect();
tokio::time::sleep(std::time::Duration::from_millis(300)).await;
cleanup.commit().await?;
for turn in turns {
turn.await??;
}
assert_eq!(
conversation_and_memory_counts(&db, conv_id).await?.0,
1,
"the turns must have created the conversation again, once"
);
assert_eq!(
count(
&db,
"SELECT count(*) FROM flow_conversation_message WHERE conversation_id = $1",
conv_id,
)
.await?,
2,
"both turns' messages should be there"
);
Ok(())
}
#[sqlx::test(fixtures("base"))]
async fn test_clear_schedule_removes_side_rows(db: Pool<Postgres>) -> anyhow::Result<()> {
initialize_tracing().await;
@@ -192,6 +370,74 @@ async fn test_workspace_delete_removes_side_rows(db: Pool<Postgres>) -> anyhow::
Ok(())
}
/// The purge endpoint carries its own copy of the emptied-conversation rule, so it gets the
/// same guard: the conversation and its memory go with the last message, and not before.
#[sqlx::test(fixtures("base"))]
async fn test_jobs_export_delete_removes_a_conversation_once_its_last_message_goes(
db: Pool<Postgres>,
) -> anyhow::Result<()> {
initialize_tracing().await;
let first_job = Uuid::new_v4();
let second_job = Uuid::new_v4();
insert_job(&db, WS, first_job).await?;
insert_job(&db, WS, second_job).await?;
let conv_id = Uuid::new_v4();
sqlx::query(
"INSERT INTO flow_conversation (id, workspace_id, flow_path, created_by)
VALUES ($1, $2, 'f/flow', 'test-user')",
)
.bind(conv_id)
.bind(WS)
.execute(&db)
.await?;
for job_id in [first_job, second_job] {
sqlx::query(
"INSERT INTO flow_conversation_message (conversation_id, message_type, content, job_id)
VALUES ($1, 'assistant', 'hi', $2)",
)
.bind(conv_id)
.bind(job_id)
.execute(&db)
.await?;
}
sqlx::query(
"INSERT INTO ai_agent_memory (workspace_id, conversation_id, step_id, messages)
VALUES ($1, $2, 'a', '[]'::jsonb)",
)
.bind(WS)
.bind(conv_id)
.execute(&db)
.await?;
let server = ApiServer::start(db.clone()).await?;
let port = server.addr.port();
let purge = |job_id: Uuid| async move {
reqwest::Client::new()
.post(format!("http://localhost:{port}/api/w/{WS}/jobs/delete"))
.header("Authorization", "Bearer SECRET_TOKEN")
.json(&[job_id])
.send()
.await
};
assert!(purge(first_job).await?.status().is_success());
assert_eq!(
conversation_and_memory_counts(&db, conv_id).await?,
(1, 1),
"a conversation with a message left must survive the purge endpoint too"
);
assert!(purge(second_job).await?.status().is_success());
assert_eq!(
conversation_and_memory_counts(&db, conv_id).await?,
(0, 0),
"the last message going should take the conversation and its memory"
);
Ok(())
}
/// The `/jobs/delete` purge endpoint must scope every side-table delete to the path
/// workspace. A `test-workspace` admin passing a job id from another workspace must not be
/// able to delete that workspace's job or side rows (the side tables no longer cascade, so
+12 -2
View File
@@ -668,10 +668,16 @@ pub async fn handle_chat_conversation_messages(
flow_path: &str,
run_query: &RunJobQuery,
user_message_raw: Option<&Box<serde_json::value::RawValue>>,
job_id: Uuid,
) -> error::Result<()> {
// Names the query parameter rather than the field: it is not a flow argument, and
// supplying it as one is the first thing tried on reading `memory_id is required`.
let memory_id = run_query.memory_id.ok_or_else(|| {
windmill_common::error::Error::BadRequest(
"memory_id is required for chat-enabled flows".to_string(),
"memory_id is required for chat-enabled flows. Pass it as the `memory_id` query \
parameter, not as a flow argument: it names the conversation the turn belongs to, \
so a fresh UUID starts one and reusing a UUID continues it."
.to_string(),
)
})?;
@@ -698,10 +704,13 @@ pub async fn handle_chat_conversation_messages(
)
.await?;
// The run this message started. Its args are the only record of what the message
// carried besides its text — attachments and every other flow input — and nothing
// written later points at them: an assistant row holds the AI agent step's job.
add_message_to_conversation_tx(
tx,
memory_id,
None,
Some(job_id),
&user_message,
MessageType::User,
None,
@@ -826,6 +835,7 @@ pub async fn run_flow<'c>(
&flow_path.to_string(),
&run_query,
args.args.get("user_message"),
uuid,
)
.await?;
}
+55 -5
View File
@@ -692,16 +692,64 @@ pub async fn delete_jobs(
.await?
.rows_affected();
let conversation_message_deleted = sqlx::query!(
// One row per message deleted, so the conversation of a chat losing several appears
// several times: the count is taken before the dedup below.
let mut conversation_ids: Vec<Uuid> = sqlx::query_scalar!(
"DELETE FROM flow_conversation_message m
USING flow_conversation c
WHERE m.conversation_id = c.id AND c.workspace_id = $1 AND m.job_id = ANY($2)",
WHERE m.conversation_id = c.id AND c.workspace_id = $1 AND m.job_id = ANY($2)
RETURNING m.conversation_id",
&w_id,
&job_ids
)
.execute(&mut *tx)
.await?
.rows_affected();
.fetch_all(&mut *tx)
.await?;
let conversation_message_deleted = conversation_ids.len() as u64;
// Same rule, lock and statement order as retention (windmill_common::jobs::delete_jobs,
// which says why): a conversation with no messages left goes, and its memory with it.
conversation_ids.sort_unstable();
conversation_ids.dedup();
let mut memory_deleted = 0;
let mut conversation_deleted = 0;
if !conversation_ids.is_empty() {
sqlx::query_scalar!(
"SELECT id FROM flow_conversation WHERE id = ANY($1) AND workspace_id = $2 ORDER BY id FOR UPDATE",
&conversation_ids,
&w_id
)
.fetch_all(&mut *tx)
.await?;
memory_deleted = sqlx::query!(
"DELETE FROM ai_agent_memory a
USING flow_conversation c
WHERE c.id = ANY($1)
AND c.workspace_id = $2
AND a.conversation_id = c.id
AND a.workspace_id = c.workspace_id
AND NOT EXISTS (
SELECT 1 FROM flow_conversation_message m WHERE m.conversation_id = c.id
)",
&conversation_ids,
&w_id
)
.execute(&mut *tx)
.await?
.rows_affected();
conversation_deleted = sqlx::query!(
"DELETE FROM flow_conversation c
WHERE c.id = ANY($1)
AND c.workspace_id = $2
AND NOT EXISTS (
SELECT 1 FROM flow_conversation_message m WHERE m.conversation_id = c.id
)",
&conversation_ids,
&w_id
)
.execute(&mut *tx)
.await?
.rows_affected();
}
// Resolutions are not exported, so a delete-then-reimport of the same UUID would
// otherwise resurrect the old annotation on a job that never carried one.
@@ -737,6 +785,8 @@ pub async fn delete_jobs(
+ zombie_deleted
+ dispatch_event_deleted
+ conversation_message_deleted
+ memory_deleted
+ conversation_deleted
+ resolution_deleted
+ jobs_deleted;
+1
View File
@@ -9553,6 +9553,7 @@ async fn run_preview_flow_job(
&flow_path,
&run_query,
user_message.as_ref(),
uuid,
)
.await?;
}
@@ -36,30 +36,20 @@ pub async fn get_or_create_conversation_with_id(
title: &str,
conversation_id: Uuid,
) -> Result<FlowConversation> {
// Check if conversation already exists
let existing_conversation = sqlx::query_as!(
FlowConversation,
"SELECT id, workspace_id, flow_path, title, created_at, updated_at, created_by
FROM flow_conversation
WHERE id = $1 AND workspace_id = $2",
conversation_id,
w_id
)
.fetch_optional(&mut **tx)
.await?;
if let Some(existing) = existing_conversation {
if let Some(existing) = lock_conversation(tx, w_id, conversation_id).await? {
return Ok(existing);
}
// Truncate title to 25 characters max
let title = truncate_with_ellipsis(title, 25);
// Create new conversation with provided ID
let conversation = sqlx::query_as!(
// Every turn released by the same collector's commit finds no row: the first insert
// wins, the others wait on it, do nothing, and read the row it created.
let created = sqlx::query_as!(
FlowConversation,
"INSERT INTO flow_conversation (id, workspace_id, flow_path, created_by, title)
VALUES ($1, $2, $3, $4, $5)
ON CONFLICT (id) DO NOTHING
RETURNING id, workspace_id, flow_path, title, created_at, updated_at, created_by",
conversation_id,
w_id,
@@ -67,10 +57,41 @@ pub async fn get_or_create_conversation_with_id(
username,
title
)
.fetch_one(&mut **tx)
.fetch_optional(&mut **tx)
.await?;
if let Some(conversation) = created {
return Ok(conversation);
}
Ok(conversation)
lock_conversation(tx, w_id, conversation_id)
.await?
.ok_or_else(|| {
crate::error::Error::BadRequest(format!(
"conversation {conversation_id} belongs to another workspace"
))
})
}
/// Locked, so a turn orders against retention collecting the conversation
/// (windmill_common::jobs::delete_jobs): either the turn goes first and the collector then
/// sees its message, or it waits and finds the row gone and creates it again. Unlocked, the
/// message insert would wait on the parent row's lock instead and then fail its FK check.
async fn lock_conversation(
tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
w_id: &str,
conversation_id: Uuid,
) -> Result<Option<FlowConversation>> {
Ok(sqlx::query_as!(
FlowConversation,
"SELECT id, workspace_id, flow_path, title, created_at, updated_at, created_by
FROM flow_conversation
WHERE id = $1 AND workspace_id = $2
FOR UPDATE",
conversation_id,
w_id
)
.fetch_optional(&mut **tx)
.await?)
}
/// Add a message to a conversation using an existing transaction
+52 -3
View File
@@ -478,6 +478,12 @@ pub static WORKER_INTERNAL_SERVER_INLINE_UTILS: OnceCell<WorkerInternalServerInl
/// set-based deletes below cost one scan per table per call instead. Because the cascade no
/// longer fires, every code path that deletes from `v2_job` by id must go through this helper
/// (or delete these tables itself) or it will leave orphan rows behind.
/// **Transaction contract:** call this inside a transaction. The conversation cleanup below
/// locks rows to serialise itself against a concurrent delete, and on an autocommit
/// connection that lock is released at statement end, silently restoring the race.
/// A conversation is collected only once every message row of it has gone with a job; a
/// row written with no job id (an MCP tool call, persisted under no job of its own) keeps
/// its conversation and the agent's memory for it alive for as long as it exists.
pub async fn delete_jobs(conn: &mut sqlx::PgConnection, ids: &[uuid::Uuid]) -> error::Result<()> {
sqlx::query!(
"DELETE FROM dispatch_event WHERE producer_job_id = ANY($1)",
@@ -485,12 +491,55 @@ pub async fn delete_jobs(conn: &mut sqlx::PgConnection, ids: &[uuid::Uuid]) -> e
)
.execute(&mut *conn)
.await?;
sqlx::query!(
"DELETE FROM flow_conversation_message WHERE job_id = ANY($1)",
let mut conversation_ids: Vec<uuid::Uuid> = sqlx::query_scalar!(
"DELETE FROM flow_conversation_message WHERE job_id = ANY($1) RETURNING conversation_id",
ids
)
.execute(&mut *conn)
.fetch_all(&mut *conn)
.await?;
conversation_ids.sort_unstable();
conversation_ids.dedup();
if !conversation_ids.is_empty() {
// A conversation is a view over its messages: once the last one goes with its job,
// the row and the agent's memory for it are all that is left, and nothing else
// collects them — `ai_agent_memory` carries no job id for retention to match on.
// Two statements rather than one CTE: a data-modifying CTE reads the snapshot from
// before the delete above, so every conversation would still look non-empty.
// Two calls each deleting one of a conversation's last messages would each still see
// the other's row — uncommitted deletes are invisible across transactions — so
// neither would collect it and nothing would try again. Taking the conversation row
// first serialises them: the second reads the first's delete and finds it empty.
sqlx::query_scalar!(
"SELECT id FROM flow_conversation WHERE id = ANY($1) ORDER BY id FOR UPDATE",
&conversation_ids
)
.fetch_all(&mut *conn)
.await?;
// Memory first, since it reads the conversation row for its workspace.
sqlx::query!(
"DELETE FROM ai_agent_memory a
USING flow_conversation c
WHERE c.id = ANY($1)
AND a.conversation_id = c.id
AND a.workspace_id = c.workspace_id
AND NOT EXISTS (
SELECT 1 FROM flow_conversation_message m WHERE m.conversation_id = c.id
)",
&conversation_ids
)
.execute(&mut *conn)
.await?;
sqlx::query!(
"DELETE FROM flow_conversation c
WHERE c.id = ANY($1)
AND NOT EXISTS (
SELECT 1 FROM flow_conversation_message m WHERE m.conversation_id = c.id
)",
&conversation_ids
)
.execute(&mut *conn)
.await?;
}
sqlx::query!("DELETE FROM zombie_job_counter WHERE job_id = ANY($1)", ids)
.execute(&mut *conn)
.await?;