use meerkat_core::service::CreateSessionRequest;
#[derive(Debug)]
pub struct RecoveredCreateRequest {
pub request: CreateSessionRequest,
pub runtime_was_registered: bool,
}
#[derive(Clone, Copy, Debug)]
pub enum RecoveryRuntimeBindingMode {
Authoritative,
LocalResources,
}
#[must_use]
pub fn unknown_provider_message(provider: &str) -> String {
format!("unknown provider '{provider}' (expected anthropic, openai, gemini, or self_hosted)")
}
pub fn parse_provider_override(provider: &str) -> Result<meerkat_core::Provider, String> {
meerkat_core::Provider::parse_strict(provider).ok_or_else(|| unknown_provider_message(provider))
}
#[cfg(all(feature = "session-store", not(target_arch = "wasm32")))]
mod context {
use std::sync::Arc;
use meerkat_core::service::SessionError;
use meerkat_core::types::SessionId;
use meerkat_core::{
AgentLlmClientDecorator, AgentToolDispatcher, ConfigRuntime, RuntimeBuildMode, Session,
SurfaceSessionRecoveryContext, build_recovered_session, connection::RealmId,
};
use meerkat_runtime::MeerkatMachine;
use super::{RecoveredCreateRequest, RecoveryRuntimeBindingMode};
use crate::factory::encode_llm_client_override_for_service;
use crate::service_factory::FactoryAgentBuilder;
use crate::session_runtime::errors::RecoveryError;
use meerkat_session::PersistentSessionService;
pub struct RecoveryContext<'a> {
pub service: &'a Arc<PersistentSessionService<FactoryAgentBuilder>>,
pub runtime_adapter: &'a Arc<MeerkatMachine>,
pub realm_id: Option<&'a RealmId>,
pub instance_id: Option<&'a str>,
pub backend: Option<&'a str>,
pub default_llm_client: Option<Arc<dyn meerkat_client::LlmClient>>,
pub agent_llm_client_decorator: Option<AgentLlmClientDecorator>,
pub external_tools: Option<Arc<dyn AgentToolDispatcher>>,
pub config_runtime: Option<Arc<ConfigRuntime>>,
}
impl RecoveryContext<'_> {
pub async fn load_persisted_session(
&self,
session_id: &SessionId,
) -> Result<Option<Session>, SessionError> {
let Some(session) = self.service.load_authoritative_session(session_id).await? else {
return Ok(None);
};
if self
.service
.session_archived_by_authority(session_id, &session)
.await?
{
return Ok(None);
}
Ok(Some(session))
}
pub async fn recovered_create_request(
&self,
session_id: &SessionId,
session: Session,
overrides: meerkat_core::SurfaceSessionRecoveryOverrides,
) -> Result<RecoveredCreateRequest, RecoveryError> {
self.recovered_create_request_with_runtime_binding_mode(
session_id,
session,
overrides,
RecoveryRuntimeBindingMode::Authoritative,
)
.await
}
pub async fn recovered_create_request_with_runtime_binding_mode(
&self,
session_id: &SessionId,
session: Session,
overrides: meerkat_core::SurfaceSessionRecoveryOverrides,
binding_mode: RecoveryRuntimeBindingMode,
) -> Result<RecoveredCreateRequest, RecoveryError> {
let current_generation = match self.config_runtime.as_ref() {
Some(runtime) => runtime.get().await.ok().map(|snapshot| snapshot.generation),
None => None,
};
let runtime_was_registered = self.runtime_adapter.contains_session(session_id).await;
let bindings = match binding_mode {
RecoveryRuntimeBindingMode::Authoritative => {
self.runtime_adapter
.prepare_bindings(session_id.clone())
.await
}
RecoveryRuntimeBindingMode::LocalResources => {
self.runtime_adapter
.prepare_local_session_bindings(session_id.clone())
.await
}
}
.map_err(|e| RecoveryError::BindingPreparation {
session_id: session_id.clone(),
message: e.to_string(),
})?;
let recovered = match build_recovered_session(
session,
&overrides,
SurfaceSessionRecoveryContext {
llm_client_override: self
.default_llm_client
.as_ref()
.map(|client| encode_llm_client_override_for_service(Arc::clone(client))),
agent_llm_client_decorator: self.agent_llm_client_decorator.clone(),
external_tools: self.external_tools.clone(),
checkpointer: None,
runtime_build_mode: RuntimeBuildMode::SessionOwned(bindings),
realm_id: self.realm_id.cloned(),
instance_id: self.instance_id.map(ToString::to_string),
backend: self.backend.map(ToString::to_string),
config_generation: current_generation,
},
) {
Ok(recovered) => recovered,
Err(error) => {
if !runtime_was_registered
&& let Err(cleanup_error) =
self.runtime_adapter.unregister_session(session_id).await
{
return Err(RecoveryError::BindingPreparation {
session_id: session_id.clone(),
message: format!(
"{error}; additionally failed to unregister newly recovered runtime binding: {cleanup_error}"
),
});
}
return Err(RecoveryError::Recovery(error));
}
};
Ok(RecoveredCreateRequest {
request: recovered.into_deferred_create_request(),
runtime_was_registered,
})
}
}
}
#[cfg(all(feature = "session-store", not(target_arch = "wasm32")))]
pub use context::RecoveryContext;