Keep messages in sync (#7850)
Co-authored-by: Douwe Osinga <douwe@squareup.com>
This commit is contained in:
@@ -1180,7 +1180,6 @@ impl Agent {
|
|||||||
).await?;
|
).await?;
|
||||||
|
|
||||||
let mut no_tools_called = true;
|
let mut no_tools_called = true;
|
||||||
let mut messages_to_add = Conversation::default();
|
|
||||||
let mut tools_updated = false;
|
let mut tools_updated = false;
|
||||||
let mut did_recovery_compact_this_iteration = false;
|
let mut did_recovery_compact_this_iteration = false;
|
||||||
let mut exit_chat = false;
|
let mut exit_chat = false;
|
||||||
@@ -1235,7 +1234,8 @@ impl Agent {
|
|||||||
if !text.is_empty() {
|
if !text.is_empty() {
|
||||||
last_assistant_text = text;
|
last_assistant_text = text;
|
||||||
}
|
}
|
||||||
messages_to_add.push(response.clone());
|
session_manager.add_message(&session_config.id, &response).await?;
|
||||||
|
conversation.push(response);
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1431,7 +1431,8 @@ impl Agent {
|
|||||||
response.created,
|
response.created,
|
||||||
thinking_content,
|
thinking_content,
|
||||||
).with_id(format!("msg_{}", Uuid::new_v4()));
|
).with_id(format!("msg_{}", Uuid::new_v4()));
|
||||||
messages_to_add.push(thinking_msg);
|
session_manager.add_message(&session_config.id, &thinking_msg).await?;
|
||||||
|
conversation.push(thinking_msg);
|
||||||
}
|
}
|
||||||
|
|
||||||
// Collect reasoning content to attach to tool request messages
|
// Collect reasoning content to attach to tool request messages
|
||||||
@@ -1459,11 +1460,14 @@ impl Agent {
|
|||||||
request.metadata.as_ref(),
|
request.metadata.as_ref(),
|
||||||
request.tool_meta.clone(),
|
request.tool_meta.clone(),
|
||||||
);
|
);
|
||||||
messages_to_add.push(request_msg);
|
|
||||||
let final_response = tool_response_messages[idx]
|
let final_response = tool_response_messages[idx]
|
||||||
.lock().await.clone();
|
.lock().await.clone();
|
||||||
yield AgentEvent::Message(final_response.clone());
|
// Persist the tool request and response as a pair
|
||||||
messages_to_add.push(final_response);
|
session_manager.add_message(&session_config.id, &request_msg).await?;
|
||||||
|
session_manager.add_message(&session_config.id, &final_response).await?;
|
||||||
|
conversation.push(request_msg);
|
||||||
|
conversation.push(final_response.clone());
|
||||||
|
yield AgentEvent::Message(final_response);
|
||||||
} else {
|
} else {
|
||||||
error!(
|
error!(
|
||||||
"Tool call could not be parsed: {}",
|
"Tool call could not be parsed: {}",
|
||||||
@@ -1600,11 +1604,13 @@ impl Agent {
|
|||||||
if final_output_tool.final_output.is_none() {
|
if final_output_tool.final_output.is_none() {
|
||||||
warn!("Final output tool has not been called yet. Continuing agent loop.");
|
warn!("Final output tool has not been called yet. Continuing agent loop.");
|
||||||
let message = Message::user().with_text(FINAL_OUTPUT_CONTINUATION_MESSAGE);
|
let message = Message::user().with_text(FINAL_OUTPUT_CONTINUATION_MESSAGE);
|
||||||
messages_to_add.push(message.clone());
|
session_manager.add_message(&session_config.id, &message).await?;
|
||||||
|
conversation.push(message.clone());
|
||||||
yield AgentEvent::Message(message);
|
yield AgentEvent::Message(message);
|
||||||
} else {
|
} else {
|
||||||
let message = Message::assistant().with_text(final_output_tool.final_output.clone().unwrap());
|
let message = Message::assistant().with_text(final_output_tool.final_output.clone().unwrap());
|
||||||
messages_to_add.push(message.clone());
|
session_manager.add_message(&session_config.id, &message).await?;
|
||||||
|
conversation.push(message.clone());
|
||||||
yield AgentEvent::Message(message);
|
yield AgentEvent::Message(message);
|
||||||
exit_chat = true;
|
exit_chat = true;
|
||||||
}
|
}
|
||||||
@@ -1615,6 +1621,8 @@ impl Agent {
|
|||||||
Ok(should_retry) => {
|
Ok(should_retry) => {
|
||||||
if should_retry {
|
if should_retry {
|
||||||
info!("Retry logic triggered, restarting agent loop");
|
info!("Retry logic triggered, restarting agent loop");
|
||||||
|
session_manager.replace_conversation(&session_config.id, &conversation).await?;
|
||||||
|
yield AgentEvent::HistoryReplaced(conversation.clone());
|
||||||
} else {
|
} else {
|
||||||
exit_chat = true;
|
exit_chat = true;
|
||||||
}
|
}
|
||||||
@@ -1655,17 +1663,14 @@ impl Agent {
|
|||||||
}).await?;
|
}).await?;
|
||||||
}
|
}
|
||||||
conversation = Conversation::new_unvalidated(updated_messages);
|
conversation = Conversation::new_unvalidated(updated_messages);
|
||||||
messages_to_add.push(summary_msg);
|
session_manager.add_message(&session_config.id, &summary_msg).await?;
|
||||||
|
conversation.push(summary_msg);
|
||||||
} else {
|
} else {
|
||||||
warn!("Expected a tool request/reply pair, but found {} matching messages",
|
warn!("Expected a tool request/reply pair, but found {} matching messages",
|
||||||
matching.len());
|
matching.len());
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
for msg in &messages_to_add {
|
|
||||||
session_manager.add_message(&session_config.id, msg).await?;
|
|
||||||
}
|
|
||||||
conversation.extend(messages_to_add);
|
|
||||||
if exit_chat {
|
if exit_chat {
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user