use crate::client::InternalClient;
use crate::retry_await;
use kortex_gen_grpc::hstp::v1::{upsert_response, CollisionStrategy, ErrorCode};
use crate::types::{entity::HSMLEntity, error::HstpError};
pub async fn upsert(
client: &mut InternalClient,
entities: &[HSMLEntity],
collision_strategy: CollisionStrategy,
retries: u32,
) -> Result<Vec<HSMLEntity>, HstpError> {
let requests = entities
.iter()
.map(|entity| {
let entity_string: Result<String, HstpError> = entity.try_into();
let entity_string = match entity_string {
Ok(entity_string) => entity_string,
Err(e) => return Err(e),
};
Ok(kortex_gen_grpc::hstp::v1::UpsertRequest {
collision_strategy: collision_strategy as i32,
entity: entity_string,
})
})
.collect::<Result<Vec<_>, HstpError>>()?;
let batch_upsert_request = kortex_gen_grpc::hstp::v1::BatchUpsertRequest { requests };
let response = retry_await!(retries, client.batch_upsert(batch_upsert_request.clone()))
.map_err(|e| {
HstpError::new(
ErrorCode::None,
format!("Failed to upsert entities: {}", e),
"".into(),
)
})?
.into_inner();
response
.responses
.into_iter()
.map(|response| match response.response {
Some(upsert_response::Response::Entity(response)) => HSMLEntity::try_from(&response),
Some(upsert_response::Response::Error(error)) => Err(HstpError::from(error)),
None => Err(HstpError::new(
ErrorCode::None,
"Failed to upsert entities, one of the responses was none",
"".into(),
)),
})
.collect::<Result<Vec<_>, HstpError>>()
}