use std::sync::Arc;
use anyhow::{Context, Result};
use theway_core::ThinkingLevel;
use theway_core::multiagent::jobs::JobTranscriptStore;
use crate::hook_executors::daemon_executors;
use crate::hooks;
use crate::runtime_storage::{RuntimeStorage, SessionRepository};
use crate::triggers;
#[derive(Clone)]
pub struct SessionProjectResources {
pub memory_block: String,
pub skills: Vec<theway_core::Skill>,
pub templates: Vec<theway_core::PromptTemplate>,
pub memory_dir: std::path::PathBuf,
pub reload_skills_fn: theway_core::ReloadSkillsFn,
pub load_local_sources: bool,
}
impl SessionProjectResources {
pub async fn load(
paths: &crate::DaemonPaths,
cli_builtin_skills: &[String],
config_builtin_skills: &[String],
load_local_sources: bool,
) -> Result<Self> {
let memory_dir = paths.base.join("memory");
let memory_block = crate::tools::memory::load_memory_block(&memory_dir).await;
let loaded_skills = if load_local_sources {
crate::skills::load_all(paths).await
} else {
crate::skills::LoadedSkills {
skills: Vec::new(),
diagnostics: Vec::new(),
}
};
let loaded_templates = if load_local_sources {
crate::templates::load_all(paths).await
} else {
crate::templates::LoadedTemplates {
templates: Vec::new(),
diagnostics: Vec::new(),
}
};
let resolved_builtins =
crate::builtin_skills::resolve_builtins(cli_builtin_skills, config_builtin_skills)?;
let mut skills = crate::builtin_skills::merge_with_user_project(
resolved_builtins.skills.clone(),
&loaded_skills.skills,
);
let state = if load_local_sources {
crate::skill_overrides::load(&paths.base).await
} else {
crate::skill_overrides::SkillOverrides::default()
};
crate::skill_overrides::apply(&state, &mut skills);
let reload_skills_fn: theway_core::ReloadSkillsFn = {
let paths = paths.clone();
let builtins = resolved_builtins.skills.clone();
std::sync::Arc::new(move || {
let paths = paths.clone();
let builtins = builtins.clone();
Box::pin(async move {
let loaded = if load_local_sources {
crate::skills::load_all(&paths).await
} else {
crate::skills::LoadedSkills {
skills: Vec::new(),
diagnostics: Vec::new(),
}
};
let mut merged =
crate::builtin_skills::merge_with_user_project(builtins, &loaded.skills);
let state = if load_local_sources {
crate::skill_overrides::load(&paths.base).await
} else {
crate::skill_overrides::SkillOverrides::default()
};
crate::skill_overrides::apply(&state, &mut merged);
theway_core::LoadSkillsOutput {
skills: merged,
diagnostics: loaded.diagnostics,
}
})
})
};
Ok(Self {
memory_block,
skills,
templates: loaded_templates.templates,
memory_dir,
reload_skills_fn,
load_local_sources,
})
}
}
#[derive(Clone, Default)]
pub struct SessionMcpResources {
pub tools: Vec<Arc<dyn theway_core::AgentTool>>,
pub notification_hooks: Arc<parking_lot::Mutex<Vec<Arc<triggers::McpNotificationHook>>>>,
pub inject_summary_servers: std::collections::HashSet<String>,
pub inject_and_run_servers: std::collections::HashSet<String>,
pub server_count: usize,
pub server_names: Vec<String>,
pub tool_names: Vec<String>,
pub notification_hook_count: usize,
}
impl SessionMcpResources {
pub fn from_loaded(loaded: crate::mcp_loader::LoadedMcp) -> Self {
for diagnostic in &loaded.diagnostics {
tracing::warn!(target: "mcp", "{diagnostic}");
}
let tool_names = loaded
.tools
.iter()
.map(|tool| tool.definition().name.clone())
.collect::<Vec<_>>();
let notification_hook_count = loaded.notification_hooks.len();
Self {
tools: loaded.tools,
notification_hooks: Arc::new(parking_lot::Mutex::new(loaded.notification_hooks)),
inject_summary_servers: loaded.inject_summary_servers,
inject_and_run_servers: loaded.inject_and_run_servers,
server_count: loaded.client_count,
server_names: loaded.server_names,
tool_names,
notification_hook_count,
}
}
}
#[derive(Clone)]
pub struct SessionExtensionResources {
pub compact_algorithms:
std::sync::Arc<theway_core::agent::compaction::algorithm::CompactAlgorithmRegistry>,
pub legacy_compaction_host: Option<std::sync::Arc<crate::ts_extensions::LegacyCompactionHost>>,
pub runtime_extension_packages:
std::sync::Arc<parking_lot::RwLock<crate::ts_extensions::PackageCatalog>>,
pub runtime_extension_engine: Option<std::sync::Arc<crate::ts_extensions::QuickJsEnginePool>>,
}
impl SessionExtensionResources {
pub fn new(
cwd: &std::path::Path,
base: &std::path::Path,
executor: std::sync::Arc<dyn theway_core::executor::ToolExecutor>,
load_local_sources: bool,
) -> Self {
let ts_extensions = if load_local_sources {
for warning in theway_extensions::ensure_managed_installed(base) {
tracing::warn!(target: "extensions", "{warning}");
}
crate::ts_extensions::ExtensionRegistry::discover(cwd, base)
} else {
crate::ts_extensions::ExtensionRegistry::new()
};
for error in &ts_extensions.errors {
tracing::warn!(target: "extensions", "{error}");
}
let legacy_compaction_host = std::sync::Arc::new(
crate::ts_extensions::LegacyCompactionHost::new(&ts_extensions),
);
let compact_algorithms = legacy_compaction_host.registry();
let runtime_extension_packages = std::sync::Arc::new(parking_lot::RwLock::new(
ts_extensions.package_catalog().clone(),
));
let runtime_extension_engine = load_local_sources.then(|| {
let broker_services =
crate::ts_extensions::ExtensionBrokerServices::new(base, executor);
for package in runtime_extension_packages.read().effective_packages() {
for permission in package.granted_permissions() {
if let theway_contract::extension::ExtensionPermission::SecretsRead(name) =
permission
&& let Ok(value) = std::env::var(name)
{
broker_services.set_secret(name, value);
}
}
}
std::sync::Arc::new(
crate::ts_extensions::QuickJsEnginePool::with_broker_services(
std::thread::available_parallelism()
.map(usize::from)
.unwrap_or(1)
.min(4),
crate::ts_extensions::QuickJsEngineLimits::default(),
broker_services,
),
)
});
Self {
compact_algorithms,
legacy_compaction_host: Some(legacy_compaction_host),
runtime_extension_packages,
runtime_extension_engine,
}
}
}
#[derive(Clone)]
pub struct SessionHookResources {
pub(super) loaded: Arc<hooks::LoadedHooks>,
}
impl SessionHookResources {
pub async fn load(paths: &crate::DaemonPaths, read_local_files: bool) -> Self {
let loaded = hooks::load_with(
paths,
"",
None::<&theway_llm_provider::Model>,
None::<ThinkingLevel>,
daemon_executors(),
read_local_files,
)
.await;
for diag in &loaded.diagnostics {
tracing::warn!(target: "hooks", "hooks loader: {diag}");
}
Self {
loaded: Arc::new(loaded),
}
}
pub fn loaded_hooks(
&self,
session_id: impl Into<String>,
model: Option<&theway_llm_provider::Model>,
thinking_level: Option<ThinkingLevel>,
) -> hooks::LoadedHooks {
hooks::LoadedHooks {
runner: Arc::new(
self.loaded
.runner
.for_session(session_id, model, thinking_level),
),
diagnostics: self.loaded.diagnostics.clone(),
}
}
}
#[derive(Clone)]
pub struct SessionExecutionContext {
pub session_id: String,
pub cwd: std::path::PathBuf,
pub transcript_store: Arc<dyn JobTranscriptStore>,
pub repo: Arc<dyn SessionRepository>,
pub storage: Arc<dyn RuntimeStorage>,
pub paths: crate::DaemonPaths,
pub executor: Arc<dyn theway_core::executor::ToolExecutor>,
pub model: Option<theway_llm_provider::Model>,
pub thinking: theway_core::ThinkingLevel,
pub resources: SessionProjectResources,
pub mcp: SessionMcpResources,
pub hooks: SessionHookResources,
pub extension_resources: SessionExtensionResources,
}
impl SessionExecutionContext {
pub fn new(
session_id: impl Into<String>,
cwd: std::path::PathBuf,
repo: Arc<dyn SessionRepository>,
storage: Arc<dyn RuntimeStorage>,
paths: crate::DaemonPaths,
executor: Arc<dyn theway_core::executor::ToolExecutor>,
model: impl Into<Option<theway_llm_provider::Model>>,
thinking: theway_core::ThinkingLevel,
resources: SessionProjectResources,
mcp: SessionMcpResources,
hooks: SessionHookResources,
) -> Self {
let session_id = session_id.into();
let cwd = cwd.canonicalize().unwrap_or(cwd);
let transcript_store = storage.job_transcript_store(&cwd);
let paths = paths.with_work_dir(cwd.clone());
let extension_resources = SessionExtensionResources::new(
&cwd,
&paths.base,
executor.clone(),
resources.load_local_sources,
);
Self {
session_id,
cwd,
transcript_store,
repo,
storage,
paths,
executor,
model: model.into(),
thinking,
resources,
mcp,
hooks,
extension_resources,
}
}
#[allow(dead_code)] pub async fn build_for_work_dir(
session_id: impl Into<String>,
requested_work_dir: std::path::PathBuf,
repo: Arc<dyn SessionRepository>,
storage: Arc<dyn RuntimeStorage>,
base_paths: crate::DaemonPaths,
model: impl Into<Option<theway_llm_provider::Model>>,
thinking: theway_core::ThinkingLevel,
cli_builtin_skills: &[String],
config_builtin_skills: &[String],
load_local_sources: bool,
) -> Result<Self> {
let cwd = requested_work_dir
.canonicalize()
.with_context(|| format!("canonicalize work dir {}", requested_work_dir.display()))?;
let paths = base_paths.with_work_dir(cwd.clone());
let executor = crate::executor::executor_for_cwd(cwd.clone());
let loaded_mcp = if load_local_sources {
crate::mcp_loader::load_all(&paths).await
} else {
crate::mcp_loader::LoadedMcp::empty()
};
let resources = SessionProjectResources::load(
&paths,
cli_builtin_skills,
config_builtin_skills,
load_local_sources,
)
.await?;
let hooks = SessionHookResources::load(&paths, load_local_sources).await;
Ok(SessionExecutionContext::new(
session_id,
cwd,
repo,
storage,
base_paths,
executor,
model,
thinking,
resources,
SessionMcpResources::from_loaded(loaded_mcp),
hooks,
))
}
}