use std::sync::Arc;
use zeph_core::agent::Agent;
use zeph_core::channel::LoopbackChannel;
use super::deps::ServeAgentDeps;
const REACTIVATION_LOCK_RETRY_ATTEMPTS: u32 = 5;
const REACTIVATION_LOCK_RETRY_DELAY: std::time::Duration = std::time::Duration::from_millis(20);
#[tracing::instrument(
name = "serve.agent_factory.build",
skip_all,
level = "info",
fields(session_id = session_id.as_str())
)]
pub(crate) async fn build_agent_factory(
deps: ServeAgentDeps,
session_id: zeph_common::SessionId,
conversation_id: zeph_memory::ConversationId,
) -> impl FnOnce(LoopbackChannel) -> Agent<LoopbackChannel> + Send + 'static {
let (session_sink, preloaded_messages) = if deps.session_persistence_config.enabled {
hydrate_session_sink(
&deps.session_persistence_config,
&deps.memory,
session_id.clone(),
conversation_id,
&deps.resume_condenser,
deps.resume_token_counter.as_ref(),
deps.session_config.budget_tokens,
)
.await
} else {
(None, Vec::new())
};
move |channel| {
let debug_config = deps.session_config.debug_config.clone();
let mut agent = Agent::new_with_registry_arc(
deps.provider,
deps.embedding_provider,
channel,
deps.registry,
deps.matcher,
deps.max_active_skills,
zeph_tools::DynExecutor(deps.tool_executor),
)
.apply_session_config(deps.session_config)
.with_skill_matching_config(
deps.skill_disambiguation_threshold,
deps.skill_two_stage_matching,
deps.skill_confusability_threshold,
)
.with_skill_provider_names(
deps.skill_generation_provider,
deps.skill_disambiguate_provider,
)
.with_semantic_scan(deps.semantic_scan, deps.semantic_scan_provider)
.with_memory(
Arc::clone(&deps.memory),
conversation_id,
deps.history_limit,
deps.recall_limit,
deps.summarization_threshold,
)
.with_session_sink(session_sink)
.with_session_persistence_config(Some(deps.session_persistence_config))
.with_provider_pool(deps.provider_pool, deps.provider_config_snapshot);
if !preloaded_messages.is_empty() {
agent = agent.with_preloaded_messages(preloaded_messages);
}
if debug_config.enabled {
let session_dump_dir = debug_config.output_dir.join(session_id.as_str());
agent = crate::agent_setup::apply_debug_dumper(
agent,
session_dump_dir.as_path(),
debug_config.format,
)
.0;
}
agent
}
}
async fn hydrate_session_sink(
session_persistence_config: &zeph_config::SessionConfig,
memory: &Arc<zeph_memory::semantic::SemanticMemory>,
session_id: zeph_common::SessionId,
conversation_id: zeph_memory::ConversationId,
resume_condenser: &zeph_session::LlmCondenser,
resume_token_counter: &zeph_agent_context::memory_backend::TokenCounterAdapter,
context_window: usize,
) -> (
Option<Arc<zeph_agent_persistence::SessionSink>>,
Vec<zeph_llm::provider::Message>,
) {
let data_dir = std::path::PathBuf::from(&session_persistence_config.data_dir);
let session_path = zeph_session::session_dir(&data_dir, session_id.as_str());
let store = zeph_session::SessionStore::new(memory.sqlite().pool().clone());
if let Err(e) = store.create(session_id.as_str()).await {
tracing::warn!(error = %e, "serve-sessions: failed to seed session metadata row");
}
if let Err(e) = store
.link_conversation(session_id.as_str(), conversation_id.0)
.await
{
tracing::warn!(error = %e, "serve-sessions: failed to link session to conversation");
}
for attempt in 1..=REACTIVATION_LOCK_RETRY_ATTEMPTS {
match zeph_agent_persistence::hydrate_and_condense(
&session_path,
&store,
session_id.as_str(),
conversation_id,
memory,
None,
resume_condenser,
resume_token_counter,
context_window,
)
.await
{
Ok(hydrated) => {
let sink = Arc::new(zeph_agent_persistence::SessionSink::new(
hydrated.log,
store,
session_id,
));
return (Some(sink), hydrated.messages);
}
Err(zeph_agent_persistence::PersistenceError::Session(
zeph_session::SessionError::AlreadyLocked(lock_path),
)) if attempt < REACTIVATION_LOCK_RETRY_ATTEMPTS => {
tracing::debug!(
lock_path,
attempt,
"serve-sessions: event log still locked by a draining prior actor; retrying reactivation"
);
tokio::time::sleep(REACTIVATION_LOCK_RETRY_DELAY).await;
}
Err(e) => {
if matches!(
e,
zeph_agent_persistence::PersistenceError::Session(
zeph_session::SessionError::AlreadyLocked(_)
)
) {
metrics::counter!("serve.session.reactivation_lock_exhausted_total")
.increment(1);
}
tracing::warn!(error = %e, "serve-sessions: session persistence disabled for this session");
return (None, Vec::new());
}
}
}
unreachable!(
"the loop always returns: the last attempt's AlreadyLocked falls into the generic Err(e) arm"
)
}
#[cfg(test)]
mod tests {
use super::*;
use zeph_agent_context::memory_backend::TokenCounterAdapter;
use zeph_llm::any::AnyProvider;
use zeph_memory::semantic::SemanticMemory;
async fn make_memory() -> Arc<SemanticMemory> {
Arc::new(
SemanticMemory::new(
":memory:",
"http://127.0.0.1:1",
None,
AnyProvider::Mock(zeph_llm::mock::MockProvider::default()),
"test-model",
)
.await
.unwrap(),
)
}
fn make_test_condenser() -> (zeph_session::LlmCondenser, TokenCounterAdapter) {
let deps = zeph_context::summarization::SummarizationDeps {
provider: AnyProvider::Mock(zeph_llm::mock::MockProvider::default()),
llm_timeout: std::time::Duration::from_secs(5),
token_counter: std::sync::Arc::new(TokenCounterAdapter::new(std::sync::Arc::new(
zeph_memory::TokenCounter::new(),
))),
structured_summaries: true,
on_anchored_summary: None,
};
let condenser = zeph_session::LlmCondenser::new(deps, 1.0, 1);
let token_counter_adapter =
TokenCounterAdapter::new(std::sync::Arc::new(zeph_memory::TokenCounter::new()));
(condenser, token_counter_adapter)
}
#[tokio::test]
async fn hydrate_session_sink_links_conversation_id() {
let memory = make_memory().await;
let cid = memory.sqlite().create_conversation().await.unwrap();
let dir = tempfile::tempdir().unwrap();
let config = zeph_config::SessionConfig {
enabled: true,
data_dir: dir.path().to_string_lossy().into_owned(),
..Default::default()
};
let session_id = zeph_common::SessionId::new("s1");
let (condenser, token_counter) = make_test_condenser();
let (sink, messages) = hydrate_session_sink(
&config,
&memory,
session_id.clone(),
cid,
&condenser,
&token_counter,
0,
)
.await;
assert!(
sink.is_some(),
"session persistence enabled must produce a SessionSink"
);
assert!(
messages.is_empty(),
"a brand-new session has no history to replay"
);
let store = zeph_session::SessionStore::new(memory.sqlite().pool().clone());
let meta = store.get(session_id.as_str()).await.unwrap().unwrap();
assert_eq!(
meta.conversation_id,
Some(cid.0),
"hydrate_session_sink must link the session to its conversation_id, or reactivation \
can never find one to replay"
);
}
#[tokio::test]
async fn hydrate_session_sink_replays_prior_history_on_reactivation() {
let memory = make_memory().await;
let cid = memory.sqlite().create_conversation().await.unwrap();
let dir = tempfile::tempdir().unwrap();
let config = zeph_config::SessionConfig {
enabled: true,
data_dir: dir.path().to_string_lossy().into_owned(),
..Default::default()
};
let session_id = zeph_common::SessionId::new("s1");
let (condenser, token_counter) = make_test_condenser();
let (sink, _) = hydrate_session_sink(
&config,
&memory,
session_id.clone(),
cid,
&condenser,
&token_counter,
0,
)
.await;
let sink = sink.expect("session persistence enabled must produce a SessionSink");
sink.record_message(zeph_llm::provider::Role::User, "hello", &[])
.await
.unwrap();
drop(sink);
let (_, messages) = hydrate_session_sink(
&config,
&memory,
session_id,
cid,
&condenser,
&token_counter,
0,
)
.await;
assert_eq!(
messages.len(),
1,
"reactivation must replay the session's prior turn, not start fresh"
);
}
#[tokio::test]
async fn hydrate_session_sink_retries_transient_already_locked_and_succeeds() {
let memory = make_memory().await;
let cid = memory.sqlite().create_conversation().await.unwrap();
let dir = tempfile::tempdir().unwrap();
let config = zeph_config::SessionConfig {
enabled: true,
data_dir: dir.path().to_string_lossy().into_owned(),
..Default::default()
};
let session_id = zeph_common::SessionId::new("s1");
let session_path = zeph_session::session_dir(dir.path(), session_id.as_str());
let (condenser, token_counter) = make_test_condenser();
let blocker = zeph_session::SessionEventLog::open_exclusive(&session_path)
.await
.unwrap();
let release_handle = tokio::spawn(async move {
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
drop(blocker);
});
let (sink, _) = hydrate_session_sink(
&config,
&memory,
session_id,
cid,
&condenser,
&token_counter,
0,
)
.await;
release_handle.await.unwrap();
assert!(
sink.is_some(),
"hydrate_session_sink must retry past a transient AlreadyLocked and eventually \
succeed once the draining actor releases its flock"
);
}
#[tokio::test]
async fn hydrate_session_sink_gives_up_after_retry_budget_exhausted() {
let memory = make_memory().await;
let cid = memory.sqlite().create_conversation().await.unwrap();
let dir = tempfile::tempdir().unwrap();
let config = zeph_config::SessionConfig {
enabled: true,
data_dir: dir.path().to_string_lossy().into_owned(),
..Default::default()
};
let session_id = zeph_common::SessionId::new("s1");
let session_path = zeph_session::session_dir(dir.path(), session_id.as_str());
let (condenser, token_counter) = make_test_condenser();
let _blocker = zeph_session::SessionEventLog::open_exclusive(&session_path)
.await
.unwrap();
let (sink, messages) = hydrate_session_sink(
&config,
&memory,
session_id,
cid,
&condenser,
&token_counter,
0,
)
.await;
assert!(
sink.is_none(),
"retry budget exhaustion against a genuinely still-locked session must degrade to \
no persistence, not panic or hang"
);
assert!(messages.is_empty());
}
#[tokio::test]
async fn build_agent_factory_wires_debug_dumper_when_enabled() {
let memory = make_memory().await;
let cid = memory.sqlite().create_conversation().await.unwrap();
let dump_dir = tempfile::tempdir().unwrap();
let session_id = zeph_common::SessionId::new("debug-test-session");
let mut config = zeph_core::config::Config::default();
config.debug.enabled = true;
config.debug.output_dir = dump_dir.path().to_path_buf();
config.debug.format = zeph_config::DumpFormat::Raw;
let session_config = zeph_core::AgentSessionConfig::from_config(&config, 4096);
let (condenser, token_counter) = make_test_condenser();
let deps = ServeAgentDeps {
provider: AnyProvider::Mock(zeph_llm::mock::MockProvider::default()),
embedding_provider: AnyProvider::Mock(zeph_llm::mock::MockProvider::default()),
registry: Arc::new(parking_lot::RwLock::new(
zeph_skills::registry::SkillRegistry::empty(),
)),
matcher: None,
max_active_skills: 5,
skill_disambiguation_threshold: 0.2,
skill_two_stage_matching: false,
skill_confusability_threshold: 0.0,
skill_generation_provider: String::new(),
skill_disambiguate_provider: String::new(),
semantic_scan: false,
semantic_scan_provider: String::new(),
tool_executor: Arc::new(zeph_tools::SetCwdExecutor),
memory: Arc::clone(&memory),
history_limit: 10,
recall_limit: 5,
summarization_threshold: 1000,
session_config,
session_persistence_config: zeph_config::SessionConfig::default(),
resume_condenser: Arc::new(condenser),
resume_token_counter: Arc::new(token_counter),
provider_pool: Vec::new(),
provider_config_snapshot: zeph_core::ProviderConfigSnapshot::default(),
};
let build_agent = build_agent_factory(deps, session_id.clone(), cid).await;
let (channel, _handle) = zeph_core::LoopbackChannel::pair(8);
let _agent = build_agent(channel);
let session_dump_dir = dump_dir.path().join(session_id.as_str());
assert!(
session_dump_dir.is_dir(),
"debug dump session subdirectory must be created when [debug] enabled = true"
);
let has_timestamped_child = std::fs::read_dir(&session_dump_dir)
.unwrap()
.next()
.is_some();
assert!(
has_timestamped_child,
"DebugDumper::new must create a timestamped subdirectory under the session dir"
);
}
#[tokio::test]
async fn build_agent_factory_wires_provider_pool_for_background_resolution() {
use zeph_commands::traits::agent::AgentAccess as _;
let memory = make_memory().await;
let cid = memory.sqlite().create_conversation().await.unwrap();
let session_id = zeph_common::SessionId::new("provider-pool-test-session");
let config = zeph_core::config::Config::default();
let session_config = zeph_core::AgentSessionConfig::from_config(&config, 4096);
let (condenser, token_counter) = make_test_condenser();
let deps = ServeAgentDeps {
provider: AnyProvider::Mock(zeph_llm::mock::MockProvider::default()),
embedding_provider: AnyProvider::Mock(zeph_llm::mock::MockProvider::default()),
registry: Arc::new(parking_lot::RwLock::new(
zeph_skills::registry::SkillRegistry::empty(),
)),
matcher: None,
max_active_skills: 5,
skill_disambiguation_threshold: 0.2,
skill_two_stage_matching: false,
skill_confusability_threshold: 0.0,
skill_generation_provider: String::new(),
skill_disambiguate_provider: String::new(),
semantic_scan: false,
semantic_scan_provider: String::new(),
tool_executor: Arc::new(zeph_tools::SetCwdExecutor),
memory: Arc::clone(&memory),
history_limit: 10,
recall_limit: 5,
summarization_threshold: 1000,
session_config,
session_persistence_config: zeph_config::SessionConfig::default(),
resume_condenser: Arc::new(condenser),
resume_token_counter: Arc::new(token_counter),
provider_pool: vec![zeph_core::config::ProviderEntry {
name: Some("named-test".into()),
model: Some("llama3.2".into()),
..zeph_core::config::ProviderEntry::default()
}],
provider_config_snapshot: zeph_core::ProviderConfigSnapshot::default(),
};
let build_agent = build_agent_factory(deps, session_id.clone(), cid).await;
let (channel, _handle) = zeph_core::LoopbackChannel::pair(8);
let mut agent = build_agent(channel);
let output = agent.handle_provider("").await;
assert!(
output.contains("named-test"),
"the pool entry configured on ServeAgentDeps must be visible through the built \
Agent's provider_pool; got: {output}"
);
}
fn embed_fn_constant(text: &str) -> zeph_skills::matcher::EmbedFuture {
let _ = text;
Box::pin(async { Ok(vec![1.0_f32, 0.0]) })
}
#[tokio::test]
async fn build_agent_factory_wires_skill_matching_config() {
use zeph_commands::traits::agent::AgentAccess as _;
let memory = make_memory().await;
let cid = memory.sqlite().create_conversation().await.unwrap();
let session_id = zeph_common::SessionId::new("skill-matching-config-test-session");
let config = zeph_core::config::Config::default();
let session_config = zeph_core::AgentSessionConfig::from_config(&config, 4096);
let (condenser, token_counter) = make_test_condenser();
let skill_meta = zeph_skills::loader::SkillMeta {
name: "solo-skill".to_owned(),
description: "a lone skill with no confusable sibling".to_owned(),
..Default::default()
};
let inner_matcher =
zeph_skills::matcher::SkillMatcher::new(&[&skill_meta], embed_fn_constant)
.await
.expect("single-skill matcher construction must succeed with a constant embed_fn");
let deps = ServeAgentDeps {
provider: AnyProvider::Mock(zeph_llm::mock::MockProvider::default()),
embedding_provider: AnyProvider::Mock(zeph_llm::mock::MockProvider::default()),
registry: Arc::new(parking_lot::RwLock::new(
zeph_skills::registry::SkillRegistry::empty(),
)),
matcher: Some(zeph_skills::matcher::SkillMatcherBackend::InMemory(
inner_matcher,
)),
max_active_skills: 5,
skill_disambiguation_threshold: 0.77,
skill_two_stage_matching: true,
skill_confusability_threshold: 0.42,
skill_generation_provider: "gen".to_owned(),
skill_disambiguate_provider: "dis".to_owned(),
semantic_scan: false,
semantic_scan_provider: String::new(),
tool_executor: Arc::new(zeph_tools::SetCwdExecutor),
memory: Arc::clone(&memory),
history_limit: 10,
recall_limit: 5,
summarization_threshold: 1000,
session_config,
session_persistence_config: zeph_config::SessionConfig::default(),
resume_condenser: Arc::new(condenser),
resume_token_counter: Arc::new(token_counter),
provider_pool: Vec::new(),
provider_config_snapshot: zeph_core::ProviderConfigSnapshot::default(),
};
let build_agent = build_agent_factory(deps, session_id.clone(), cid).await;
let (channel, _handle) = zeph_core::LoopbackChannel::pair(8);
let mut agent = build_agent(channel);
let output = agent
.handle_skills("confusability")
.await
.expect("handle_skills(\"confusability\") must not error");
assert!(
output.contains("above 0.42"),
"ServeAgentDeps::skill_confusability_threshold = 0.42 must reach the built Agent's \
ConfusabilityReport exactly (not e.g. 0.77, disambiguation_threshold's value, from a \
swapped with_skill_matching_config argument); got: {output}"
);
}
#[tokio::test]
async fn build_agent_factory_wires_semantic_scan_config() {
use zeph_commands::traits::agent::AgentAccess as _;
let memory = make_memory().await;
let cid = memory.sqlite().create_conversation().await.unwrap();
let session_id = zeph_common::SessionId::new("semantic-scan-config-test-session");
let config = zeph_core::config::Config::default();
let session_config = zeph_core::AgentSessionConfig::from_config(&config, 4096);
let (condenser, token_counter) = make_test_condenser();
let deps = ServeAgentDeps {
provider: AnyProvider::Mock(zeph_llm::mock::MockProvider::default()),
embedding_provider: AnyProvider::Mock(zeph_llm::mock::MockProvider::default()),
registry: Arc::new(parking_lot::RwLock::new(
zeph_skills::registry::SkillRegistry::empty(),
)),
matcher: None,
max_active_skills: 5,
skill_disambiguation_threshold: 0.2,
skill_two_stage_matching: false,
skill_confusability_threshold: 0.0,
skill_generation_provider: String::new(),
skill_disambiguate_provider: String::new(),
semantic_scan: true,
semantic_scan_provider: "scan-test-provider".to_owned(),
tool_executor: Arc::new(zeph_tools::SetCwdExecutor),
memory: Arc::clone(&memory),
history_limit: 10,
recall_limit: 5,
summarization_threshold: 1000,
session_config,
session_persistence_config: zeph_config::SessionConfig::default(),
resume_condenser: Arc::new(condenser),
resume_token_counter: Arc::new(token_counter),
provider_pool: Vec::new(),
provider_config_snapshot: zeph_core::ProviderConfigSnapshot::default(),
};
let build_agent = build_agent_factory(deps, session_id.clone(), cid).await;
let (channel, _handle) = zeph_core::LoopbackChannel::pair(8);
let mut agent = build_agent(channel);
let err = agent
.handle_plugins("add /nonexistent/plugin/path")
.await
.expect_err(
"handle_plugins(\"add\") must fail closed when semantic_scan is enabled and \
semantic_scan_provider is not in the (empty) provider pool",
);
let msg = err.to_string();
assert!(
msg.contains("semantic_scan_provider") && msg.contains("scan-test-provider"),
"ServeAgentDeps::semantic_scan/semantic_scan_provider = true/\"scan-test-provider\" \
must reach the built Agent's fail-closed plugin-add check exactly; got: {msg}"
);
}
}