use crate::executor::error::{ExecutorError, ExecutorResult};
use crate::executor::request::{ExecutionContext, RequestContext};
use crate::storage::InOutItem;
use crate::types::io::{InputItem, ResponsesInput, resolve_tool_choice, resolve_tools};
use crate::types::request_response::RequestPayload;
use crate::utils::uuid7_str;
pub async fn rehydrate_conversation(
request: RequestPayload,
exec_ctx: &ExecutionContext,
) -> ExecutorResult<RequestContext> {
let response_id = uuid7_str("resp_");
let new_input_items: Vec<InputItem> = Vec::from(&request.input);
let original_request = request.clone();
let mut ctx = RequestContext {
enriched_request: request,
original_request,
new_input_items,
response_id,
conversation_id: None,
};
if ctx.original_request.conversation_id.is_some() && ctx.original_request.previous_response_id.is_some() {
return Err(ExecutorError::InvalidRequest(
"provide only one of conversation_id or previous_response_id".into(),
));
}
if ctx.original_request.conversation_id.is_some() {
from_conversation(&mut ctx, exec_ctx).await?;
return Ok(ctx);
}
if ctx.original_request.previous_response_id.is_some() {
from_response(&mut ctx, exec_ctx).await?;
return Ok(ctx);
}
ctx.enriched_request.input = ResponsesInput::Items(ctx.new_input_items.clone());
Ok(ctx)
}
async fn from_response(ctx: &mut RequestContext, exec_ctx: &ExecutionContext) -> ExecutorResult<()> {
let stored = exec_ctx.resp_handler.get(ctx).await?;
let history = exec_ctx.resp_handler.rehydrate(ctx).await?;
let mut items = InOutItem::into_input_items(history);
items.reserve(ctx.new_input_items.len());
items.extend(ctx.new_input_items.iter().cloned());
ctx.enriched_request.previous_response_id = None;
ctx.enriched_request.input = ResponsesInput::Items(items);
ctx.enriched_request.tools = resolve_tools(
ctx.original_request.tools.as_deref(),
stored.metadata.effective_tools.as_deref(),
ctx.original_request.tools.is_some(),
);
ctx.enriched_request.tool_choice = Some(resolve_tool_choice(
ctx.original_request.tool_choice.as_ref(),
&stored.metadata.effective_tool_choice,
ctx.original_request.tool_choice.is_some(),
));
ctx.conversation_id = stored.conversation_id;
Ok(())
}
async fn from_conversation(ctx: &mut RequestContext, exec_ctx: &ExecutionContext) -> ExecutorResult<()> {
let (conv_data, history) = tokio::try_join!(
async {
if ctx.original_request.store {
exec_ctx.conv_handler.get_or_create(ctx).await
} else {
exec_ctx.conv_handler.get(ctx).await
}
},
exec_ctx.conv_handler.rehydrate(ctx),
)?;
let mut items = InOutItem::into_input_items(history);
items.reserve(ctx.new_input_items.len());
items.extend(ctx.new_input_items.iter().cloned());
ctx.enriched_request.input = ResponsesInput::Items(items);
ctx.conversation_id = Some(conv_data.conversation_id);
Ok(())
}