use std::sync::{Arc, OnceLock};
use anyhow::{Context, Result};
use theway_contract::session::SessionStore;
use theway_core::multiagent::goal;
use theway_core::multiagent::graph::engine::DagEngine;
use theway_core::{AgentHarness, AgentHarnessOptions, ThinkingLevel};
use theway_transport::feed::FeedUpdate;
use theway_transport::inbox;
use crate::orchestration::DaemonServices;
use crate::{agent_specs, tools, triggers};
mod activation_build;
mod resources;
pub(crate) use activation_build::load_persisted_dag_runs;
#[allow(unused_imports)]
pub use resources::{
SessionExecutionContext, SessionExtensionResources, SessionHookResources, SessionMcpResources,
SessionProjectResources, parse_mcp_diagnostic,
};
pub struct SessionRuntimeBuilder {
#[allow(dead_code)] pub thinking: ThinkingLevel,
pub stream_fn: theway_core::StreamFn,
pub dag_engine: Arc<DagEngine>,
pub subagent_registry: theway_core::multiagent::jobs::SubagentJobRegistry,
pub services: DaemonServices,
pub before_tool_call: Option<theway_core::BeforeToolCallHook>,
pub control_plane_hook: Option<theway_core::OnControlPlanePromptHook>,
pub control_plane_prompt_tx: Option<
tokio::sync::mpsc::UnboundedSender<crate::control_plane_prompt::PendingControlPlanePrompt>,
>,
pub after_tool_call: Option<theway_core::AfterToolCallHook>,
pub feed_tx: tokio::sync::mpsc::UnboundedSender<(String, FeedUpdate)>,
pub main_run_tx: tokio::sync::mpsc::UnboundedSender<String>,
pub debug: bool,
pub session_cells: parking_lot::Mutex<
std::collections::HashMap<String, crate::tools::skill::SkillHarnessCell>,
>,
}
pub struct SessionRuntime {
pub session_id: String,
pub cwd: std::path::PathBuf,
pub harness: Arc<AgentHarness>,
pub trigger_executor: Arc<crate::trigger_engine::execution::TriggerExecutor>,
pub tool_names: Vec<String>,
pub hooks_active: bool,
pub extension_host: Option<Arc<crate::ts_extensions::SessionPluginHost>>,
}
#[cfg(test)]
impl SessionRuntime {
pub(crate) fn for_test(session_id: impl Into<String>, harness: Arc<AgentHarness>) -> Self {
let session_id = session_id.into();
let trigger_executor = Arc::new(crate::trigger_engine::execution::TriggerExecutor::new(
harness.agent_arc(),
harness.session().clone(),
crate::trigger_engine::runtime::TriggerRuntimeConfig::default(),
None,
None,
None,
None,
None,
None,
));
Self {
session_id: session_id.clone(),
cwd: std::env::temp_dir().join("theway-test").join(session_id),
harness,
trigger_executor,
tool_names: Vec::new(),
hooks_active: false,
extension_host: None,
}
}
}
impl SessionRuntimeBuilder {
pub async fn build(&self, ctx: &SessionExecutionContext, id: &str) -> Result<SessionRuntime> {
let store = ctx
.repo
.resume(Some(id))
.await
.with_context(|| format!("open session {id}"))?;
self.build_opened(ctx, store, true).await
}
pub async fn build_opened(
&self,
ctx: &SessionExecutionContext,
store: Arc<dyn SessionStore>,
rehydrate: bool,
) -> Result<SessionRuntime> {
let (ctx, session_id, store) = self.opened_context(ctx, store).await?;
let restored = load_persisted_dag_runs(&ctx, &session_id).await?;
let skill_harness_cell = self.install_execution_context(ctx.clone(), restored);
self.assemble_opened(ctx, store, session_id, rehydrate, skill_harness_cell)
.await
}
async fn opened_context(
&self,
ctx: &SessionExecutionContext,
store: Arc<dyn SessionStore>,
) -> Result<(Arc<SessionExecutionContext>, String, Arc<dyn SessionStore>)> {
let meta = store.get_metadata_json().await?;
let session_id = meta
.get("id")
.and_then(|v| v.as_str())
.unwrap_or("?")
.to_string();
let mut owned_ctx = ctx.clone();
owned_ctx.session_id = session_id.clone();
Ok((Arc::new(owned_ctx), session_id, store))
}
async fn assemble_opened(
&self,
ctx: Arc<SessionExecutionContext>,
store: Arc<dyn SessionStore>,
session_id: String,
rehydrate: bool,
skill_harness_cell: crate::tools::skill::SkillHarnessCell,
) -> Result<SessionRuntime> {
let extension_state_store = Arc::clone(&store);
let harness_intro = crate::session_ops::read_session_metadata(store.as_ref())
.await?
.get("harnessIntroduction")
.cloned();
let session = theway_core::Session::from_store(store);
let mut tools = tools::session_tool_set_for_cwd_with_kind(
&ctx.resources.memory_dir,
&ctx.paths.base,
&self.dag_engine,
&self.subagent_registry,
ctx.model.as_ref(),
Some(&self.stream_fn),
&skill_harness_cell,
&session_id,
ctx.executor.clone(),
&self.services,
ctx.repo.clone(),
ctx.cwd.clone(),
ctx.executor_kind,
);
if let Some(provision) = ctx.mcp.provision.as_ref() {
tools.extend(provision.read().unwrap().tools.iter().cloned());
} else {
tools.extend(ctx.mcp.tools.iter().cloned());
}
let goal_harness_cell: Arc<OnceLock<Arc<AgentHarness>>> = Arc::new(OnceLock::new());
let mut opts = AgentHarnessOptions::new(ctx.model.clone(), session.clone());
opts.observer = self.subagent_registry.observer();
opts.observation_context = theway_core::ObservationContext {
session_id: Some(session_id.clone()),
..theway_core::ObservationContext::default()
};
opts.runtime_extension_cwd = ctx.cwd.to_string_lossy().into_owned();
let session_credentials = self.services.session_execution.clone();
let session_id_for_credentials = session_id.clone();
let session_credential_resolver: theway_core::GetApiKey = Arc::new(move |provider_id| {
session_credentials
.get_credential(&session_id_for_credentials, provider_id)
.map(|secret| String::from_utf8_lossy(secret.as_bytes()).into_owned())
});
let mut runtime_extension_host = None;
if let Some(engine) = &ctx.extension_resources.runtime_extension_engine {
let base_tools = tools.clone();
let extensions =
crate::ts_extensions::SessionPluginHost::load_with_state_and_legacy(
ctx.extension_resources.runtime_extension_packages.read().clone(),
engine.as_ref().clone(),
session_id.clone(),
&ctx.cwd,
crate::ts_extensions::RuntimeExtensionHostConfig::default(),
Arc::new(
theway_core::agent::runtime_extensions::PersistentSessionExtensionStatePort::new(
extension_state_store,
),
),
ctx.extension_resources.legacy_compaction_host.clone(),
Some(ctx.extension_resources.runtime_extension_packages.clone()),
)
.await;
for diagnostic in extensions
.diagnostics()
.into_iter()
.filter(|diagnostic| diagnostic.session_id.is_some())
{
tracing::warn!(
target: "extensions",
extension_id = diagnostic.extension_id,
"{}",
diagnostic.message
);
}
tools = extensions.merge_registered_tools(tools);
let session_credential_resolver = session_credential_resolver.clone();
let credential_host = Arc::clone(&extensions);
opts.get_api_key = Some(Arc::new(move |provider_id| {
if let Some(secret) = session_credential_resolver(provider_id) {
return Some(secret);
}
credential_host.provider_api_key(provider_id)
}));
opts.runtime_extension_model_context = extensions.model_context_projection();
opts.runtime_extensions = extensions.clone();
runtime_extension_host = Some((extensions, base_tools));
} else {
opts.get_api_key = Some(session_credential_resolver);
}
let tool_names = tools
.iter()
.map(|tool| tool.definition().name.clone())
.collect::<Vec<_>>();
let context_service = crate::context::service::ContextService::new(
&ctx.cwd,
&ctx.resources.memory_block,
tool_names.clone(),
harness_intro,
);
let bundle = context_service.load(&session).await?;
opts.system_prompt = bundle.system_prompt;
opts.thinking_level = ctx.thinking;
opts.tools = tools;
opts.skills = ctx.resources.skills.clone();
opts.prompt_templates = ctx.resources.templates.clone();
opts.compact_algorithms = ctx.extension_resources.compact_algorithms.clone();
opts.stream_fn = Some(self.stream_fn.clone());
opts.reload_skills_fn = Some(ctx.resources.reload_skills_fn.clone());
opts.on_turn_end = Some(goal::stop_hook(
goal_harness_cell.clone(),
self.dag_engine.clone(),
agent_specs::launch_resolver(),
self.subagent_registry.clone(),
Some(self.stream_fn.clone()),
));
opts.turn_continuation_cap = Some(goal::MAX_CONTINUATIONS);
opts.before_tool_call = self.before_tool_call.clone();
opts.on_control_plane_prompt = match &self.control_plane_prompt_tx {
Some(tx) => Some(crate::control_plane_prompt::interactive_hook_for_session(
session_id.clone(),
tx.clone(),
)),
None => self.control_plane_hook.clone(),
};
opts.after_tool_call = self.after_tool_call.clone();
let harness = std::sync::Arc::new(AgentHarness::new(opts));
if let Some((extensions, base_tools)) = &runtime_extension_host {
let harness_ref = std::sync::Arc::downgrade(&harness);
extensions.configure_reload_tool_publisher(
base_tools.clone(),
Arc::new(move |tools| {
if let Some(harness) = harness_ref.upgrade() {
harness.replace_tools(tools);
}
}),
);
}
let (inject_summary, inject_and_run) = match ctx.mcp.provision.as_ref() {
Some(provision) => {
let slot = provision.read().unwrap();
(slot.inject_summary.clone(), slot.inject_and_run.clone())
}
None => (
ctx.mcp.inject_summary_servers.clone(),
ctx.mcp.inject_and_run_servers.clone(),
),
};
let before_trigger_action = triggers::cron_action_hook(
self.services.cron.clone(),
triggers::direct_inject_action_hook(
inject_summary,
inject_and_run,
triggers::before_trigger_action_hook(self.services.dynamic_triggers.clone()),
),
);
let trigger_executor =
std::sync::Arc::new(crate::trigger_engine::execution::TriggerExecutor::new(
harness.agent_arc(),
harness.session().clone(),
crate::trigger_engine::runtime::TriggerRuntimeConfig::default(),
None,
None,
Some(before_trigger_action),
Some(self.stream_fn.clone()),
self.before_tool_call.clone(),
self.after_tool_call.clone(),
));
let mcp_notification_hooks = match ctx.mcp.provision.as_ref() {
Some(provision) => provision.read().unwrap().hooks.clone(),
None => std::mem::take(&mut *ctx.mcp.notification_hooks.lock()),
};
register_notification_hooks(
&trigger_executor,
&mcp_notification_hooks,
&ctx.cwd,
&self.services.cron,
&self.services.dynamic_triggers,
);
let _ = skill_harness_cell.set(harness.clone());
let _ = goal_harness_cell.set(harness.clone());
let _agent_broadcast = crate::turn::listener::spawn_agent_broadcast_listener(
harness.agent().subscribe_broadcast(),
session_id.clone(),
self.feed_tx.clone(),
);
let _harness_broadcast = crate::turn::listener::spawn_harness_broadcast_listener(
harness.subscribe_session_broadcast(),
session_id.clone(),
self.feed_tx.clone(),
self.debug,
);
let _ = trigger_executor.subscribe(crate::turn::listener::trigger_listener(
session_id.clone(),
self.feed_tx.clone(),
self.debug,
));
let _ = trigger_executor.subscribe(triggers::fire_once_trigger_listener(
self.services.dynamic_triggers.clone(),
));
let _ = trigger_executor.subscribe(triggers::cron_trigger_listener(
self.services.cron.clone(),
inbox::default_inbox_path(),
));
let (hook_model, hook_thinking) = {
let state = harness.agent().state();
(state.model.clone(), state.thinking_level)
};
let loaded_hooks =
ctx.hooks
.loaded_hooks(session_id.clone(), hook_model.as_ref(), hook_thinking);
let _ = harness.agent().subscribe(loaded_hooks.runner.listener());
let _ = harness.subscribe_harness(loaded_hooks.runner.harness_listener());
let main_run_tx = self.main_run_tx.clone();
let _ = trigger_executor.subscribe(std::sync::Arc::new(
move |ev: crate::trigger_engine::event::TriggerEvent| {
if let crate::trigger_engine::event::TriggerEvent::TriggerRequestsMainRun {
trace_id,
} = ev
{
let _ = main_run_tx.send(trace_id);
}
},
));
if rehydrate {
harness
.rehydrate_from_session()
.await
.with_context(|| format!("rehydrate session {session_id}"))?;
}
harness.start_runtime_extensions().await;
Ok(SessionRuntime {
session_id,
cwd: ctx.cwd.clone(),
harness,
trigger_executor,
tool_names,
hooks_active: !loaded_hooks.runner.is_empty(),
extension_host: runtime_extension_host.map(|(host, _)| host),
})
}
}
pub(crate) use crate::trigger_engine::notification_hook::NotificationHookSink;
pub(crate) fn register_notification_hooks(
sink: &(impl NotificationHookSink + ?Sized),
mcp_notification_hooks: &[Arc<triggers::McpNotificationHook>],
cwd: &std::path::Path,
cron_registry: &triggers::cron::CronRegistry,
dynamic_trigger_registry: &triggers::dynamic::DynamicTriggerRegistry,
) {
for hook in mcp_notification_hooks {
sink.register(hook.clone());
}
sink.register(Arc::new(triggers::CronNotificationHook::new(
cron_registry.clone(),
)));
sink.register(Arc::new(triggers::DynamicTriggerCheckHook::new_for_cwd(
dynamic_trigger_registry.clone(),
cwd,
)));
}
#[cfg(test)]
tests_bridge_macro::tests_bridge!("orchestration/session");