use crate::executor::error::{ExecutorError, ExecutorResult};
use crate::executor::modes::{ConversationHandler, ResponseHandler};
use crate::executor::request::RequestContext;
use crate::types::event::ResponseStatus;
use crate::types::io::OutputItem;
use crate::types::request_response::ResponsePayload;
use tracing::error;
#[must_use]
pub(crate) fn should_persist(ctx: &RequestContext) -> bool {
ctx.original_request.store
|| ctx.original_request.previous_response_id.is_some()
|| ctx.original_request.conversation_id.is_some()
}
pub(crate) async fn persist_if_needed(
payload: ResponsePayload,
ctx: RequestContext,
conv_handler: ConversationHandler,
resp_handler: ResponseHandler,
) -> ExecutorResult<()> {
if should_persist(&ctx) {
persist_response(payload, ctx, conv_handler, resp_handler)
.await
.map_err(|source| {
error!(error = ?source, "failed to persist response");
ExecutorError::Persistence(Box::new(source))
})
} else {
Ok(())
}
}
pub async fn persist_response(
payload: ResponsePayload,
ctx: RequestContext,
conv_handler: ConversationHandler,
resp_handler: ResponseHandler,
) -> ExecutorResult<()> {
if !matches!(
payload.status.parse::<ResponseStatus>().unwrap_or_default(),
ResponseStatus::Completed | ResponseStatus::Incomplete
) || payload.id.is_empty()
{
return Ok(());
}
persist_turn(ctx, payload.output, &conv_handler, &resp_handler).await
}
pub async fn persist_turn(
ctx: RequestContext,
output_items: Vec<OutputItem>,
conv_handler: &ConversationHandler,
resp_handler: &ResponseHandler,
) -> ExecutorResult<()> {
if ctx.original_request.conversation_id.is_some() {
conv_handler.execute_turn(ctx, output_items).await
} else {
resp_handler.execute_turn(ctx, output_items).await
}
}