pub(crate) mod builder;
pub(crate) mod dispatch;
mod environment;
use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex};
use mentra::{BuiltinProvider, ModelInfo, ModelSelector, Session, agent::AgentConfig};
pub use builder::RuntimeBuilder;
use crate::{approval::SideEffectLevels, run::RunError};
use dispatch::{HookDispatch, HookRegistration, WorkspaceGuardEntry};
pub struct Runtime {
mentra: mentra::Runtime,
provider: BuiltinProvider,
provider_label: String,
model: ModelSelector,
dispatch: Arc<HookDispatch>,
levels: SideEffectLevels,
#[cfg(feature = "mcp")]
mcp_claims: Mutex<HashMap<String, PathBuf>>,
declared_claims: Mutex<HashMap<String, DeclaredClaim>>,
}
#[derive(Debug)]
struct DeclaredClaim {
root: PathBuf,
holders: usize,
}
impl std::fmt::Debug for Runtime {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Runtime")
.field("provider", &self.provider_label)
.field("model", &self.model)
.finish_non_exhaustive()
}
}
impl Runtime {
pub fn builder() -> RuntimeBuilder {
RuntimeBuilder::default()
}
pub fn mentra_runtime(&self) -> &mentra::Runtime {
&self.mentra
}
pub fn provider(&self) -> &str {
&self.provider_label
}
pub(crate) async fn resolve_model(
&self,
selector: Option<ModelSelector>,
) -> Result<ModelInfo, RunError> {
let selector = selector.unwrap_or_else(|| self.model.clone());
Ok(self.mentra.resolve_model(self.provider, selector).await?)
}
pub(crate) fn mint(
&self,
name: impl Into<String>,
model: ModelInfo,
config: AgentConfig,
_persist_identifier: &str,
) -> Result<Session, RunError> {
Ok(self
.mentra
.create_session_with_config(name, model, config)?)
}
pub(crate) fn resume_minted(&self, agent_id: &str) -> Result<Session, RunError> {
Ok(self.mentra.resume_session(agent_id)?)
}
pub(crate) fn side_effect_levels(&self) -> SideEffectLevels {
self.levels.clone()
}
pub(crate) fn interceptors(&self) -> &[Arc<dyn crate::hooks::Interceptor>] {
self.dispatch.interceptors()
}
pub(crate) fn register_workspace(&self, entry: WorkspaceGuardEntry) -> HookRegistration {
self.dispatch.register(entry)
}
#[cfg(feature = "mcp")]
pub(crate) fn claim_mcp_server(&self, name: &str, root: &Path) -> String {
let mut claims = self.mcp_claims.lock().expect("mcp claim map poisoned");
if !claims.contains_key(name) {
claims.insert(name.to_string(), root.to_path_buf());
return name.to_string();
}
let mut effective = format!("{name}-{}", root_suffix(root));
let mut attempt = 2_u32;
while claims.contains_key(&effective) {
effective = format!("{name}-{}-{attempt}", root_suffix(root));
attempt += 1;
}
claims.insert(effective.clone(), root.to_path_buf());
effective
}
#[cfg(feature = "mcp")]
pub(crate) fn release_mcp_claim(&self, name: &str, root: &Path) {
let mut claims = self.mcp_claims.lock().expect("mcp claim map poisoned");
if claims.get(name).is_some_and(|owner| owner == root) {
claims.remove(name);
}
}
pub(crate) fn claim_declared_tool(&self, name: &str, root: &Path) -> Result<(), String> {
let mut claims = self
.declared_claims
.lock()
.expect("declared tool claim map poisoned");
match claims.get_mut(name) {
Some(claim) if claim.holders > 0 && claim.root != root => Err(
"another workspace open on this runtime declares a tool by that name".to_string(),
),
Some(claim) => {
claim.root = root.to_path_buf();
claim.holders += 1;
Ok(())
}
None if self.registers_tool(name) => {
Err("this runtime already offers a tool by that name".to_string())
}
None => {
claims.insert(
name.to_string(),
DeclaredClaim {
root: root.to_path_buf(),
holders: 1,
},
);
Ok(())
}
}
}
pub(crate) fn release_declared_tool(&self, name: &str, root: &Path) {
let mut claims = self
.declared_claims
.lock()
.expect("declared tool claim map poisoned");
if let Some(claim) = claims.get_mut(name)
&& claim.root == root
{
claim.holders = claim.holders.saturating_sub(1);
}
}
pub(crate) fn foreign_declared_tools(&self, root: &Path) -> Vec<String> {
self.declared_claims
.lock()
.expect("declared tool claim map poisoned")
.iter()
.filter(|(_, claim)| claim.root != root)
.map(|(name, _)| name.clone())
.collect()
}
fn registers_tool(&self, name: &str) -> bool {
self.mentra
.tools()
.iter()
.any(|descriptor| descriptor.provider.name == name)
}
}
#[cfg(feature = "mcp")]
fn root_suffix(root: &Path) -> String {
const OFFSET: u64 = 0xcbf2_9ce4_8422_2325;
const PRIME: u64 = 0x0000_0100_0000_01b3;
let mut hash = OFFSET;
for byte in root.as_os_str().as_encoded_bytes() {
hash ^= u64::from(*byte);
hash = hash.wrapping_mul(PRIME);
}
format!("{:08x}", (hash >> 32) as u32 ^ hash as u32)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_runtime_can_be_shared_across_tasks() {
const fn assert_send_sync<T: Send + Sync>() {}
assert_send_sync::<Runtime>();
assert_send_sync::<RuntimeBuilder>();
}
#[cfg(feature = "mcp")]
#[test]
fn a_taken_server_name_is_suffixed_and_a_released_one_is_free_again() {
use std::path::Path;
let runtime = Runtime::builder()
.with_base_url("http://127.0.0.1:1/v1")
.with_api_key("test-key")
.with_ephemeral_history()
.build()
.expect("builds");
let first = runtime.claim_mcp_server("fs", Path::new("/repo/one"));
let second = runtime.claim_mcp_server("fs", Path::new("/repo/two"));
let again = runtime.claim_mcp_server("fs", Path::new("/repo/two"));
assert_eq!(first, "fs", "the first claimant keeps the plain name");
assert_ne!(second, "fs", "the second must not collide in the registry");
assert!(second.starts_with("fs-"), "{second}");
assert_ne!(again, second, "every live claim is its own namespace");
runtime.release_mcp_claim("fs", Path::new("/repo/two"));
runtime.release_mcp_claim(&second, Path::new("/repo/one"));
assert_eq!(
runtime.claim_mcp_server("fs", Path::new("/repo/three")),
format!("fs-{}", root_suffix(Path::new("/repo/three"))),
"a name someone else holds stays held"
);
runtime.release_mcp_claim("fs", Path::new("/repo/one"));
assert_eq!(
runtime.claim_mcp_server("fs", Path::new("/repo/four")),
"fs",
"a released name is claimable plain again"
);
}
}