use tracing::Instrument as _;
use super::super::{Agent, Channel};
impl<C: Channel> Agent<C> {
pub(crate) async fn check_trust_transition(&self, skill_name: &str) {
let span = tracing::info_span!("core.learning.check_trust_transition", skill = skill_name);
async move {
if let Err(_elapsed) = tokio::time::timeout(
std::time::Duration::from_secs(2),
self.check_trust_transition_inner(skill_name),
)
.await
{
tracing::warn!(
skill = skill_name,
"check_trust_transition timed out after 2s"
);
}
}
.instrument(span)
.await;
}
async fn check_trust_transition_inner(&self, skill_name: &str) {
let span = tracing::info_span!(
"core.learning.check_trust_transition_inner",
skill = skill_name
);
async move {
let Some(memory) = &self.services.memory.persistence.memory else {
return;
};
let Some(config) = &self.services.learning_engine.config else {
return;
};
let Ok(Some(metrics)) = memory.sqlite().skill_metrics(skill_name).await else {
return;
};
let successes = u32::try_from(metrics.successes).unwrap_or(0);
let failures = u32::try_from(metrics.failures).unwrap_or(0);
let total = u32::try_from(metrics.total).unwrap_or(0);
let posterior = zeph_skills::trust_score::posterior_mean(successes, failures);
if total >= config.auto_promote_min_uses && posterior > config.auto_promote_threshold {
if !cross_session_rollout_ok(
memory,
skill_name,
config.cross_session_rollout,
config.min_sessions_before_promote,
"promotion",
)
.await
{
return;
}
try_auto_promote(memory, skill_name, posterior, total).await;
}
if total >= config.auto_demote_min_uses && posterior < config.auto_demote_threshold {
if !cross_session_rollout_ok(
memory,
skill_name,
config.cross_session_rollout,
config.min_sessions_before_demote,
"demotion",
)
.await
{
return;
}
try_auto_demote(memory, skill_name, posterior, total).await;
}
}
.instrument(span)
.await;
}
}
async fn cross_session_rollout_ok(
memory: &zeph_memory::semantic::SemanticMemory,
skill_name: &str,
cross_session_rollout: bool,
min_sessions: u32,
action: &str,
) -> bool {
if !cross_session_rollout {
return true;
}
match memory.sqlite().distinct_session_count(skill_name).await {
Ok(sessions) if sessions < i64::from(min_sessions) => {
tracing::debug!(
skill = skill_name,
sessions,
required = min_sessions,
"cross-session rollout: insufficient sessions for {action}"
);
false
}
Ok(_) => true,
Err(e) => {
tracing::warn!("cross-session count query failed for {skill_name}: {e:#}");
true
}
}
}
async fn try_auto_promote(
memory: &zeph_memory::semantic::SemanticMemory,
skill_name: &str,
posterior: f64,
total: u32,
) {
let trust_level = memory
.sqlite()
.load_skill_trust(skill_name)
.await
.ok()
.flatten()
.map(|r| r.trust_level);
if trust_level != Some(zeph_common::SkillTrustLevel::Trusted)
&& trust_level != Some(zeph_common::SkillTrustLevel::Blocked)
{
tracing::info!(
skill = skill_name,
posterior = format!("{posterior:.3}"),
total,
"auto-promoting skill to trusted"
);
if trust_level.is_none() {
let _ = memory
.sqlite()
.upsert_skill_trust(
skill_name,
zeph_common::SkillTrustLevel::Trusted,
zeph_memory::store::SourceKind::Local,
None,
None,
"",
)
.await;
} else {
let _ = memory
.sqlite()
.set_skill_trust_level(skill_name, zeph_common::SkillTrustLevel::Trusted)
.await;
}
}
}
pub(crate) async fn cap_auto_version_trust(
memory: &zeph_memory::semantic::SemanticMemory,
skill_name: &str,
) -> bool {
let parent_level = match memory.sqlite().load_skill_trust(skill_name).await {
Ok(row) => row.map(|r| r.trust_level),
Err(e) => {
tracing::warn!(
skill = skill_name,
error = %e,
"cap_auto_version_trust: trust read failed, skipping activation (fail-closed)"
);
return false;
}
};
let capped = match parent_level {
Some(level) => level.min_trust(zeph_common::SkillTrustLevel::Quarantined),
None => zeph_common::SkillTrustLevel::Quarantined,
};
let trust_write_result = if parent_level.is_some() {
memory
.sqlite()
.set_skill_trust_level(skill_name, capped)
.await
} else {
memory
.sqlite()
.upsert_skill_trust(
skill_name,
capped,
zeph_memory::store::SourceKind::Local,
None,
None,
"",
)
.await
.map(|()| true)
};
match trust_write_result {
Ok(true) => {
tracing::info!(
skill = skill_name,
trust = %capped,
"store_improved_version: capped auto-improved version trust"
);
true
}
Ok(false) => {
tracing::warn!(
skill = skill_name,
"store_improved_version: trust write affected no rows, skipping activation (fail-closed)"
);
false
}
Err(e) => {
tracing::warn!(
skill = skill_name,
error = %e,
"store_improved_version: trust write failed, skipping activation (fail-closed)"
);
false
}
}
}
pub(crate) async fn scan_and_cap_for_activation(
memory: &zeph_memory::semantic::SemanticMemory,
skill_name: &str,
body: &str,
) -> bool {
let scan = zeph_skills::scanner::scan_skill_body(body);
if scan.has_matches() {
tracing::warn!(
skill = skill_name,
patterns = ?scan.matched_patterns,
"scan_and_cap_for_activation: hard-blocking activation, injection scan matched"
);
return false;
}
cap_auto_version_trust(memory, skill_name).await
}
async fn try_auto_demote(
memory: &zeph_memory::semantic::SemanticMemory,
skill_name: &str,
posterior: f64,
total: u32,
) {
let Ok(Some(trust_row)) = memory.sqlite().load_skill_trust(skill_name).await else {
return;
};
if trust_row.trust_level == zeph_common::SkillTrustLevel::Trusted
|| trust_row.trust_level == zeph_common::SkillTrustLevel::Verified
{
tracing::warn!(
skill = skill_name,
posterior = format!("{posterior:.3}"),
total,
"auto-demoting skill to quarantined"
);
let _ = memory
.sqlite()
.set_skill_trust_level(skill_name, zeph_common::SkillTrustLevel::Quarantined)
.await;
}
}