use std::collections::HashMap;
use std::path::PathBuf;
use std::sync::Arc;
use car_connectors::ConnectorStatus;
use car_eventlog::{Event, EventKind, EventLog, RetentionPolicy};
use car_memgine::maintenance::{decide_maintenance, MaintenanceDecision, MaintenanceInput};
use car_memgine::self_evolution::{
run_evolution_cycle, ComponentState, EvolutionCycleReport, EvolutionOutcome, EvolutionPolicy,
EvolutionSignals, EvolvableComponent,
};
use car_memgine::{MemgineEngine, TraceEvent};
use serde::{Deserialize, Serialize};
use serde_json::Value;
use crate::session::ServerState;
#[async_trait::async_trait]
pub trait HarnessMeasurer: Send + Sync {
async fn measure(
&self,
request: &HarnessMeasureRequest,
harness_config: Option<&car_memgine::HarnessConfig>,
memgine_config: Option<&car_memgine::MemgineConfig>,
) -> Result<car_eventlog::harness_metrics::HarnessMetrics, String>;
}
fn default_split() -> String {
"held-out".to_string()
}
fn default_held_in_fraction() -> f64 {
0.5
}
fn default_max_turns() -> u32 {
20
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct HarnessMeasureRequest {
pub model: String,
#[serde(default = "default_split")]
pub split: String,
#[serde(default = "default_held_in_fraction")]
pub held_in_fraction: f64,
#[serde(default)]
pub split_seed: u64,
#[serde(default = "default_max_turns")]
pub max_turns: u32,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub tasks_dir: Option<PathBuf>,
}
pub fn parse_harness_measure_request(
params: &Value,
) -> Result<Option<HarnessMeasureRequest>, String> {
match params.get("harness_measure") {
Some(v) if !v.is_null() => Ok(Some(
serde_json::from_value(v.clone())
.map_err(|e| format!("invalid harness_measure: {e}"))?,
)),
_ => Ok(None),
}
}
pub async fn measure_baseline(
measurer: &dyn HarnessMeasurer,
request: &HarnessMeasureRequest,
live_config: Option<&car_memgine::HarnessConfig>,
) -> Result<car_eventlog::harness_metrics::HarnessMetrics, String> {
measurer
.measure(request, live_config, None)
.await
.map_err(|e| format!("baseline harness measurement failed: {e}"))
}
pub async fn measure_candidate(
measurer: &dyn HarnessMeasurer,
request: &HarnessMeasureRequest,
base: &car_memgine::HarnessConfig,
patch: &car_memgine::harness_evolution::HarnessConfigPatch,
) -> Result<car_eventlog::harness_metrics::HarnessMetrics, String> {
let candidate = base.with_patch_for_measurement(patch);
measurer.measure(request, Some(&candidate), None).await
}
pub fn parse_context_measure_request(
params: &Value,
) -> Result<Option<HarnessMeasureRequest>, String> {
match params.get("context_measure") {
Some(v) if !v.is_null() => Ok(Some(
serde_json::from_value(v.clone())
.map_err(|e| format!("invalid context_measure: {e}"))?,
)),
_ => Ok(None),
}
}
pub async fn measure_context_baseline(
measurer: &dyn HarnessMeasurer,
request: &HarnessMeasureRequest,
live: &car_memgine::MemgineConfig,
) -> Result<car_eventlog::harness_metrics::HarnessMetrics, String> {
measurer
.measure(request, None, Some(live))
.await
.map_err(|e| format!("baseline context measurement failed: {e}"))
}
pub async fn measure_context_candidate(
measurer: &dyn HarnessMeasurer,
request: &HarnessMeasureRequest,
base: &car_memgine::MemgineConfig,
patch: &car_memgine::ContextConfigPatch,
) -> Result<car_eventlog::harness_metrics::HarnessMetrics, String> {
let candidate = base.with_context_patch_for_measurement(patch);
measurer.measure(request, None, Some(&candidate)).await
}
pub fn harness_component_from_events(events: &[Event]) -> Option<ComponentState> {
if events.is_empty() {
return None;
}
let report = car_eventlog::harness_adapt::diagnose(events, 2);
let implicated: usize = report.interventions.iter().map(|i| i.evidence_count).sum();
let pressure = (implicated as f64 / events.len() as f64).min(1.0);
Some(ComponentState {
component: EvolvableComponent::Harness,
signals: EvolutionSignals {
pressure,
evidence: events.len() as u64,
min_evidence: 20,
cost: 2.0,
},
})
}
pub fn harness_component_from_metrics(
m: &car_eventlog::harness_metrics::HarnessMetrics,
) -> Option<ComponentState> {
let eff = &m.trajectory_efficiency;
let attempts = eff.actions_succeeded + eff.failed_attempts;
if attempts == 0 {
return None;
}
Some(ComponentState {
component: EvolvableComponent::Harness,
signals: EvolutionSignals {
pressure: (eff.failed_attempts as f64 / attempts as f64).clamp(0.0, 1.0),
evidence: attempts as u64,
min_evidence: 20,
cost: 2.0,
},
})
}
pub fn tools_component_from_connectors(connectors: &[ConnectorStatus]) -> Option<ComponentState> {
if connectors.is_empty() {
return None;
}
let unhealthy = connectors.iter().filter(|c| !c.connected).count();
Some(ComponentState {
component: EvolvableComponent::Tools,
signals: EvolutionSignals {
pressure: (unhealthy as f64 / connectors.len() as f64).clamp(0.0, 1.0),
evidence: connectors.len() as u64,
min_evidence: 1,
cost: 1.0,
},
})
}
pub fn failed_trace_events(events: &[Event]) -> Vec<TraceEvent> {
events
.iter()
.filter(|ev| match ev.kind {
EventKind::ActionFailed
| EventKind::ActionRejected
| EventKind::PolicyViolation
| EventKind::ReplanExhausted => true,
EventKind::GoalEvaluated => {
let met = ev
.data
.get("met")
.and_then(|v| v.as_bool())
.unwrap_or(false);
let grounded = ev
.data
.get("grounded")
.and_then(|v| v.as_bool())
.unwrap_or(true);
met && !grounded
}
EventKind::TurnCompleted => {
let decision = ev
.data
.get("decision")
.and_then(|v| v.as_str())
.unwrap_or("");
let truncated = ev
.data
.get("was_truncated")
.and_then(|v| v.as_bool())
.unwrap_or(false);
truncated || decision == "max_turns" || decision == "stalled"
}
_ => false,
})
.map(|ev| TraceEvent {
kind: serde_json::to_value(&ev.kind)
.ok()
.and_then(|v| v.as_str().map(str::to_string))
.unwrap_or_default(),
action_id: ev.action_id.clone(),
tool: ev
.data
.get("tool")
.and_then(|v| v.as_str())
.map(str::to_string),
data: Value::Object(ev.data.clone().into_iter().collect()),
duration_ms: None,
state_before: None,
state_after: None,
reward: Some(0.0),
})
.collect()
}
pub fn maintenance_input_from_stats(stats: &car_memgine::memsys::MemoryStats) -> MaintenanceInput {
let total = stats.total_facts;
MaintenanceInput {
dirty_regions: stats.outstanding_outdated + stats.facts_superseded,
total_regions: total,
localized_cost_per_region: 1.0,
global_cost_per_region: 1.0,
global_structural_gain: if total > 0 {
(stats.facts_superseded as f64 / total as f64).clamp(0.0, 1.0)
} else {
0.0
},
gain_value: total as f64,
}
}
pub async fn run_memory_evolution(
engine: &Arc<tokio::sync::Mutex<MemgineEngine>>,
dry_run: bool,
) -> Result<EvolutionOutcome, String> {
let mut eng = engine.lock().await;
let decision: MaintenanceDecision =
decide_maintenance(&maintenance_input_from_stats(&eng.memory_stats()));
if dry_run {
return Ok(EvolutionOutcome::no_op(format!(
"dry_run: would consolidate (maintenance: {:?} — {})",
decision.strategy, decision.rationale
)));
}
let report = eng.consolidate().await;
let summary = serde_json::to_string(&serde_json::json!({
"mechanism": "consolidate",
"maintenance": decision,
"expired_pruned": report.expired_pruned,
"superseded_gc": report.superseded_gc,
"turns_compacted": report.turns_compacted,
"domains_evolved": report.domains_evolved,
"total_nodes": report.total_nodes,
}))
.map_err(|e| e.to_string())?;
Ok(EvolutionOutcome::applied(summary))
}
pub async fn run_skills_evolution(
engine: &Arc<tokio::sync::Mutex<MemgineEngine>>,
failed_events: &[TraceEvent],
dry_run: bool,
) -> Result<EvolutionOutcome, String> {
let mut eng = engine.lock().await;
if !eng.has_inference() {
return Err("no inference engine".to_string());
}
let domains = eng.domains_needing_evolution(0.6);
if domains.is_empty() {
return Ok(EvolutionOutcome::no_op(
"no domain below the evolution threshold (success < 0.6 over ≥3 outcomes) — nothing to evolve",
));
}
if dry_run {
return Ok(EvolutionOutcome::no_op(format!(
"dry_run: would evolve domain(s) {:?} over {} failure trace(s)",
domains,
failed_events.len()
)));
}
let mut evolved = 0usize;
for domain in &domains {
evolved += eng.evolve_skills(failed_events, domain).await.len();
}
let summary = format!(
"evolved {} skill(s) across domain(s) {:?} over {} failure trace(s)",
evolved,
domains,
failed_events.len()
);
Ok(if evolved > 0 {
EvolutionOutcome::applied(summary)
} else {
EvolutionOutcome::no_op(summary)
})
}
pub const TOOLS_OUT_OF_SCOPE_REASON: &str =
"connector remediation means re-running a connector's OAuth or credential exchange. That is \
an access change, and this loop deliberately holds no authority to grant, refresh or move \
credentials — reconnect and re-auth stay operator actions through the `connectors.*` \
surface. Recorded as a deliberate scope decision, not a failure.";
const CONTEXT_PENDING_REASON: &str =
"A pre-activation grade IS available for a patched context mutation: the daemon can replay \
the deterministic bench split twice — once under the live context config, once under it \
plus this patch — and promote or reject on the resulting TASK pass rates, because bench \
tasks that declare a memory fixture answer out of assembled context and are therefore \
sensitive to this knob. It is opt-in via the `context_measure` param because a benchmark \
replay is a paid side effect (real model calls, real money), and it is deliberately never \
supplied by the unattended cadence. Activation therefore falls back to the human gate here. \
Approving this fingerprint once makes every later cycle that proposes the same change apply \
it, measure the conversation tokens it actually saved over compacting without it, and roll \
it back if it saved none — that post-apply margin check is retained on the approved path as \
defence in depth, and is a weaker claim than the task-outcome grade above because it can \
only falsify the predicted token saving, never confirm the assembled context still answers. \
The approval ledger is DAEMON-WIDE and the fingerprint names the change, not the engine: \
approving it authorizes this same conversation_keep_recent change on any engine this daemon \
evolves — the shared one the unattended cadence runs over, and every per-agent engine an \
`evolution.run` names — not only the one that proposed it.";
const CONTEXT_NOT_REQUESTED: &str =
"no pre-activation grade was requested for this cycle: the `context_measure` param was \
absent, and the unattended cadence never supplies it (a timer must not start spending \
benchmark replays because someone enabled `evolution_interval_secs`).";
const CONTEXT_DRY_RUN_REASON: &str =
"`context_measure` was requested with `dry_run` — a benchmark replay is a paid side effect \
and a dry run performs none, so nothing was measured and nothing could be graded.";
fn context_patch_base_moved(
measured_under: &car_memgine::MemgineConfig,
current: &car_memgine::MemgineConfig,
) -> bool {
measured_under.conversation_keep_recent != current.conversation_keep_recent
}
fn context_config_moved_reason(measured_under: usize, current: usize) -> String {
format!(
"the live context config moved while this mutation was being measured: \
conversation_keep_recent was {measured_under} when the baseline replay ran and is \
{current} now. Something else moved it — another session's `evolution.run` over the \
same engine, the human-approved path, or an unattended cadence tick — so the grade \
was computed against a base that no longer exists, and the inverse patch this apply \
would hand back for rollback would name the CURRENT value rather than the measured \
one. Nothing is applied. The mutation is NOT falsified and no backoff is recorded: \
the measurement was invalidated, not the change, and a later cycle re-diagnoses \
against the new base and re-measures against it."
)
}
const CONTEXT_NO_PATCH_REASON: &str =
"this mutation carries no concrete config patch, so there is nothing to project into a \
candidate config and nothing to install on a replay — a human designs this change.";
fn context_pending_reason(missing: &str) -> String {
format!("{missing} {CONTEXT_PENDING_REASON}")
}
const CONTEXT_BACKOFF_REASON: &str =
"this mutation's post-apply measurement falsified it on an earlier unattended tick, so it is \
in exponential backoff. Re-applying it every tick would re-run a full compaction pass under \
the engine lock to reach the same verdict — the Skills arm backs off for the same reason. \
The standing approval is untouched: the next attempt happens automatically once the window \
elapses, and a re-diagnosis that PAYS clears the backoff.";
#[derive(Debug, Default)]
pub struct ContextBackoff {
map: HashMap<String, DomainAttempts>,
}
impl ContextBackoff {
pub fn is_due(&self, fingerprint: &str, tick: u64) -> bool {
self.map
.get(fingerprint)
.map(|a| tick >= a.next_tick)
.unwrap_or(true)
}
pub fn note_falsified(&mut self, fingerprint: &str, tick: u64) {
let entry = self
.map
.entry(fingerprint.to_string())
.or_insert(DomainAttempts {
attempts: 0,
next_tick: tick,
});
entry.attempts = (entry.attempts + 1).min(BACKOFF_MAX_EXPONENT);
entry.next_tick = tick + (1u64 << entry.attempts);
}
pub fn clear(&mut self, fingerprint: &str) {
self.map.remove(fingerprint);
}
}
fn backoff_due(backoff: Option<(&std::sync::Mutex<ContextBackoff>, u64)>, fp: &str) -> bool {
match backoff {
Some((b, tick)) => b.lock().unwrap().is_due(fp, tick),
None => true,
}
}
fn note_falsified(backoff: Option<(&std::sync::Mutex<ContextBackoff>, u64)>, fp: &str) {
if let Some((b, tick)) = backoff {
b.lock().unwrap().note_falsified(fp, tick);
}
}
fn clear_backoff(backoff: Option<(&std::sync::Mutex<ContextBackoff>, u64)>, fp: &str) {
if let Some((b, _)) = backoff {
b.lock().unwrap().clear(fp);
}
}
pub async fn run_context_evolution(
engine: &Arc<tokio::sync::Mutex<MemgineEngine>>,
state: &Arc<ServerState>,
dry_run: bool,
pending: &std::sync::Mutex<Vec<Value>>,
backoff: Option<(&std::sync::Mutex<ContextBackoff>, u64)>,
measure: Option<(&dyn HarnessMeasurer, &HarnessMeasureRequest)>,
) -> Result<EvolutionOutcome, String> {
use car_memgine::context_evolution::{
context_mutation_fingerprint, diagnose_context, requires_human_approval,
};
use car_memgine::harness_evolution::{EvolutionAgent, PromotionDecision};
let signals = { engine.lock().await.context_evolution_signals() };
let Some(signals) = signals else {
return Ok(EvolutionOutcome::no_op("no observable context signal"));
};
let mutations = diagnose_context(&signals);
if mutations.is_empty() {
return Ok(EvolutionOutcome::no_op(
"no context mutations diagnosed from live context signals",
));
}
let mut details: Vec<Value> = Vec::new();
let mut applied = 0usize;
let mut pending_count = 0usize;
let mut grade_attempts = 0usize;
for m in &mutations {
let fingerprint = context_mutation_fingerprint(m);
let prior = {
let ledger = state.approval_ledger.read().await;
ledger.lookup(&fingerprint).map(|r| r.decision)
};
let status: Value = match prior {
Some(car_policy::ApprovalDecision::Rejected) => {
serde_json::json!({ "status": "rejected_by_operator" })
}
Some(car_policy::ApprovalDecision::Approved) => match m.patch.as_ref() {
None => serde_json::json!({
"status": "approved_no_patch",
"note": "approved but carries no concrete config patch — a human designs this change",
}),
Some(_) if !backoff_due(backoff, &fingerprint) => serde_json::json!({
"status": "in_backoff",
"governance": "human_approved",
"reason": CONTEXT_BACKOFF_REASON,
}),
Some(_) if dry_run => {
serde_json::json!({ "status": "would_apply", "governance": "human_approved" })
}
Some(patch) => {
let mut eng = engine.lock().await;
let baseline = eng.compact_conversation_heuristic();
let before = eng
.context_evolution_signals()
.map(|s| s.conversation_tokens)
.unwrap_or(0);
match eng.apply_context_patch(patch) {
Err(e) => {
note_falsified(backoff, &fingerprint);
serde_json::json!({
"status": "apply_failed",
"error": e,
"baseline_turns_summarized": baseline.turns_summarized,
})
}
Ok(inverse) => {
let report = eng.compact_conversation_heuristic();
let after = eng
.context_evolution_signals()
.map(|s| s.conversation_tokens)
.unwrap_or(before);
if after >= before {
let why = if report.turns_summarized == 0 {
"compaction under the new conversation_keep_recent did no \
work at all: after the baseline pass the turns still held \
verbatim are below the engine's own hard threshold, so it \
refuses to compact them whatever this knob says. That is \
\"no saving available right now\", not \"this knob cannot \
help\" — the baseline pass has left the layer in the shape \
where a later cycle's attempt can pay, and the backoff \
window is when it retries."
} else {
"compaction under the new conversation_keep_recent ran and \
still saved nothing over compaction under the old one, so \
the change's predicted improvement is falsified."
};
let falsified = format!(
"{why} ({before} → {after} tokens.) The config is reverted; \
the summarization the measurement itself performed is not \
undone — compaction replaces turns with summaries and keeps \
them in the layer, which is what the engine's own heuristic \
does at this same saturation.",
);
note_falsified(backoff, &fingerprint);
match eng.apply_context_patch(&inverse) {
Ok(_) => serde_json::json!({
"status": "rolled_back",
"governance": "human_approved",
"reason": falsified,
"conversation_tokens_baseline": before,
"conversation_tokens_after": after,
"baseline_turns_summarized": baseline.turns_summarized,
"turns_summarized": report.turns_summarized,
}),
Err(rollback_error) => {
tracing::error!(
fingerprint = %fingerprint,
error = %rollback_error,
"context mutation was falsified but its rollback \
failed — the engine is running at the mutated \
conversation_keep_recent"
);
serde_json::json!({
"status": "rollback_failed",
"governance": "human_approved",
"reason": falsified,
"rollback_error": rollback_error,
"rollback_patch": inverse,
"conversation_tokens_baseline": before,
"conversation_tokens_after": after,
"baseline_turns_summarized": baseline.turns_summarized,
"turns_summarized": report.turns_summarized,
})
}
}
} else {
applied += 1;
clear_backoff(backoff, &fingerprint);
serde_json::json!({
"status": "applied",
"governance": "human_approved",
"rollback_patch": inverse,
"conversation_tokens_baseline": before,
"conversation_tokens_after": after,
"baseline_turns_summarized": baseline.turns_summarized,
"turns_summarized": report.turns_summarized,
})
}
}
}
}
},
None => {
let gradeable = m.patch.as_ref().filter(|_| !requires_human_approval(m));
match (measure, gradeable) {
(Some((measurer, request)), Some(patch)) if !dry_run => {
let live = { engine.lock().await.config().clone() };
grade_attempts += 1;
let replays = match measure_context_baseline(measurer, request, &live).await
{
Err(e) => Err(e),
Ok(baseline) => {
match measure_context_candidate(measurer, request, &live, patch)
.await
{
Err(e) => Err(e),
Ok(candidate) => Ok((baseline, candidate)),
}
}
};
match replays {
Err(error) => serde_json::json!({
"status": "measurement_failed",
"error": error,
}),
Ok((baseline, candidate)) => {
let decision = EvolutionAgent::new()
.evaluate_context(m, &baseline, &candidate);
let mut status = match decision {
PromotionDecision::Promote { reason } => {
let mut eng = engine.lock().await;
let current = eng.config().clone();
if context_patch_base_moved(&live, ¤t) {
serde_json::json!({
"status": "config_moved_during_measurement",
"governance": "promoted",
"reason": context_config_moved_reason(
live.conversation_keep_recent,
current.conversation_keep_recent,
),
"gate_reason": reason,
"measured_under_conversation_keep_recent":
live.conversation_keep_recent,
"current_conversation_keep_recent":
current.conversation_keep_recent,
})
} else {
match eng.apply_context_patch(patch) {
Ok(inverse) => {
applied += 1;
serde_json::json!({
"status": "applied",
"governance": "promoted",
"reason": reason,
"rollback_patch": inverse,
})
}
Err(e) => serde_json::json!({
"status": "apply_failed",
"governance": "promoted",
"error": e,
}),
}
}
}
PromotionDecision::Reject { reason } => serde_json::json!({
"status": "rejected_by_gate",
"reason": reason,
}),
PromotionDecision::NeedsApproval { reason }
| PromotionDecision::Incomparable { reason } => {
let reason = context_pending_reason(&reason);
pending_count += 1;
pending.lock().unwrap().push(serde_json::json!({
"fingerprint": fingerprint,
"mutation": m.id,
"component": m.contract.component,
"safety_affecting": m.contract.component.is_safety_affecting(),
"rationale": m.rationale,
"reason": reason,
}));
serde_json::json!({
"status": "pending_approval",
"reason": reason,
})
}
};
if let Some(obj) = status.as_object_mut() {
obj.insert(
"baseline_task_pass_rate".into(),
serde_json::to_value(baseline.task_pass_rate)
.unwrap_or(Value::Null),
);
obj.insert(
"baseline_task_pass_denominator".into(),
serde_json::to_value(baseline.task_pass_denominator)
.unwrap_or(Value::Null),
);
obj.insert(
"baseline_total_tokens".into(),
Value::from(baseline.trajectory_efficiency.total_tokens),
);
obj.insert(
"candidate_task_pass_rate".into(),
serde_json::to_value(candidate.task_pass_rate)
.unwrap_or(Value::Null),
);
obj.insert(
"candidate_task_pass_denominator".into(),
serde_json::to_value(candidate.task_pass_denominator)
.unwrap_or(Value::Null),
);
obj.insert(
"candidate_total_tokens".into(),
Value::from(candidate.trajectory_efficiency.total_tokens),
);
}
status
}
}
}
_ => {
let missing = if gradeable.is_none() {
CONTEXT_NO_PATCH_REASON
} else if measure.is_none() {
CONTEXT_NOT_REQUESTED
} else {
CONTEXT_DRY_RUN_REASON
};
let reason = context_pending_reason(missing);
pending_count += 1;
pending.lock().unwrap().push(serde_json::json!({
"fingerprint": fingerprint,
"mutation": m.id,
"component": m.contract.component,
"safety_affecting": m.contract.component.is_safety_affecting(),
"rationale": m.rationale,
"reason": reason,
}));
serde_json::json!({
"status": "pending_approval",
"reason": reason,
})
}
}
}
};
let mut d = serde_json::json!({
"mutation": m.id,
"component": m.contract.component,
"fingerprint": fingerprint,
"rationale": m.rationale,
});
if let (Some(obj), Some(s)) = (d.as_object_mut(), status.as_object()) {
for (k, v) in s {
obj.insert(k.clone(), v.clone());
}
}
details.push(d);
}
let mut summary_obj = serde_json::json!({
"mechanism": "context_evolution",
"mutations": mutations.len(),
"applied": applied,
"pending": pending_count,
"details": details,
});
if let (Some(obj), Some((_, request))) = (summary_obj.as_object_mut(), measure) {
obj.insert(
"context_measured".into(),
serde_json::json!({
"status": if dry_run { "skipped_dry_run" } else { "measured" },
"grade_attempts": grade_attempts,
"model": request.model,
"split": request.split,
"split_seed": request.split_seed,
}),
);
}
let summary = serde_json::to_string(&summary_obj).map_err(|e| e.to_string())?;
Ok(if applied > 0 {
EvolutionOutcome::applied(summary)
} else {
EvolutionOutcome::no_op(summary)
})
}
#[derive(Debug, Default)]
pub struct SkillsBackoff {
map: HashMap<String, DomainAttempts>,
}
#[derive(Debug)]
struct DomainAttempts {
attempts: u32,
next_tick: u64,
}
const BACKOFF_MAX_EXPONENT: u32 = 6;
impl SkillsBackoff {
pub fn due(&self, flagged: &[String], tick: u64) -> Vec<String> {
flagged
.iter()
.filter(|d| {
self.map
.get(*d)
.map(|a| tick >= a.next_tick)
.unwrap_or(true)
})
.cloned()
.collect()
}
pub fn note_attempt(&mut self, domain: &str, tick: u64) {
let entry = self
.map
.entry(domain.to_string())
.or_insert(DomainAttempts {
attempts: 0,
next_tick: tick,
});
entry.attempts = (entry.attempts + 1).min(BACKOFF_MAX_EXPONENT);
entry.next_tick = tick + (1u64 << entry.attempts);
}
pub fn reset_recovered(&mut self, flagged: &[String]) {
self.map.retain(|d, _| flagged.iter().any(|f| f == d));
}
}
pub async fn run_skills_evolution_backoff(
engine: &Arc<tokio::sync::Mutex<MemgineEngine>>,
backoff: &std::sync::Mutex<SkillsBackoff>,
tick: u64,
) -> Result<EvolutionOutcome, String> {
let mut eng = engine.lock().await;
if !eng.has_inference() {
return Err("no inference engine".to_string());
}
let flagged = eng.domains_needing_evolution(0.6);
let due = {
let mut b = backoff.lock().unwrap();
b.reset_recovered(&flagged);
b.due(&flagged, tick)
};
if flagged.is_empty() {
return Ok(EvolutionOutcome::no_op(
"no domain below the evolution threshold — nothing to evolve",
));
}
if due.is_empty() {
return Ok(EvolutionOutcome::no_op(format!(
"all {} flagged domain(s) in backoff — no attempt this tick",
flagged.len()
)));
}
let mut evolved = 0usize;
for domain in &due {
evolved += eng.evolve_skills(&[], domain).await.len();
backoff.lock().unwrap().note_attempt(domain, tick);
}
let summary = format!(
"evolved {} skill(s) across due domain(s) {:?} ({} flagged total)",
evolved,
due,
flagged.len()
);
Ok(if evolved > 0 {
EvolutionOutcome::applied(summary)
} else {
EvolutionOutcome::no_op(summary)
})
}
#[derive(Debug, Default)]
pub struct CycleGuard {
running: std::sync::atomic::AtomicBool,
}
pub struct CycleToken<'a> {
guard: &'a CycleGuard,
}
impl CycleGuard {
pub fn try_begin(&self) -> Option<CycleToken<'_>> {
self.running
.compare_exchange(
false,
true,
std::sync::atomic::Ordering::SeqCst,
std::sync::atomic::Ordering::SeqCst,
)
.ok()
.map(|_| CycleToken { guard: self })
}
}
impl Drop for CycleToken<'_> {
fn drop(&mut self) {
self.guard
.running
.store(false, std::sync::atomic::Ordering::SeqCst);
}
}
pub fn seed_evolution_interval() -> Option<u64> {
let anchor = std::env::var_os("CAR_PROJECT_DIR")
.map(std::path::PathBuf::from)
.or_else(|| std::env::current_dir().ok())?;
let car_dir = car_memgine::project::discover_project(&anchor)?;
car_memgine::project::load_config_overrides(&car_dir)?
.evolution_interval_secs
.filter(|s| *s > 0)
}
pub async fn run_evolution_cadence_cycle(
state: &Arc<ServerState>,
backoff: &std::sync::Mutex<SkillsBackoff>,
context_backoff: &std::sync::Mutex<ContextBackoff>,
tick: u64,
) -> Option<EvolutionCycleReport> {
let engine = state.shared_memgine.as_ref()?.clone();
let mut components = { engine.lock().await.evolution_component_states() };
state.ensure_connectors_loaded().await;
let connector_list = state.connectors().list().await;
if let Some(t) = tools_component_from_connectors(&connector_list) {
components.push(t);
}
let context_pending: std::sync::Mutex<Vec<Value>> = std::sync::Mutex::new(Vec::new());
let context_pending_ref = &context_pending;
let policy = EvolutionPolicy::default();
let report = run_evolution_cycle(&components, &policy, |c| {
let engine = engine.clone();
async move {
match c {
EvolvableComponent::Memory => run_memory_evolution(&engine, false).await,
EvolvableComponent::Skills => {
run_skills_evolution_backoff(&engine, backoff, tick).await
}
EvolvableComponent::Harness => Ok(EvolutionOutcome::out_of_scope(
"harness telemetry and the HITL apply path are per-session — the cadence has \
no session, so harness evolution is driven via evolution.run on one. \
Recorded as a deliberate scope decision, not a failure.",
)),
EvolvableComponent::Context => {
run_context_evolution(
&engine,
state,
false,
context_pending_ref,
Some((context_backoff, tick)),
None,
)
.await
}
EvolvableComponent::Tools => {
Ok(EvolutionOutcome::out_of_scope(TOOLS_OUT_OF_SCOPE_REASON))
}
}
}
})
.await;
let pending = context_pending.into_inner().unwrap();
if !pending.is_empty() {
tracing::info!(
count = pending.len(),
"evolution cadence surfaced context mutation(s) awaiting operator approval; \
approve by fingerprint via permission.approve"
);
}
Some(report)
}
const EVOLUTION_LOG_MAX_EVENTS: usize = 1000;
pub fn spawn_evolution_cadence(
state: Arc<ServerState>,
interval_secs: u64,
) -> Option<tokio::task::JoinHandle<()>> {
if state.shared_memgine.is_none() {
tracing::warn!(
"evolution_interval_secs set but the daemon has no shared engine; cadence not started"
);
return None;
}
let journal = state.journal_dir.join("evolution.jsonl");
Some(tokio::spawn(async move {
let mut log = EventLog::with_journal(journal);
log.set_retention(Some(RetentionPolicy {
max_events: Some(EVOLUTION_LOG_MAX_EVENTS),
max_age_secs: None,
}));
let guard = Arc::new(CycleGuard::default());
let backoff = Arc::new(std::sync::Mutex::new(SkillsBackoff::default()));
let context_backoff = Arc::new(std::sync::Mutex::new(ContextBackoff::default()));
let mut tick_no: u64 = 0;
let mut tick = tokio::time::interval(std::time::Duration::from_secs(interval_secs.max(1)));
tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
tick.tick().await;
loop {
tick.tick().await;
tick_no += 1;
let Some(_token) = guard.try_begin() else {
tracing::warn!("evolution cadence tick skipped: previous cycle still running");
continue;
};
let cycle_state = state.clone();
let cycle_backoff = backoff.clone();
let cycle_context_backoff = context_backoff.clone();
let outcome = tokio::spawn(async move {
run_evolution_cadence_cycle(
&cycle_state,
&cycle_backoff,
&cycle_context_backoff,
tick_no,
)
.await
})
.await;
let mut data: HashMap<String, Value> = HashMap::new();
data.insert("source".into(), Value::from("cadence"));
match outcome {
Ok(Some(report)) => {
if report.plan.evolve_now.is_empty() && report.steps.is_empty() {
continue;
}
data.insert(
"report".into(),
serde_json::to_value(&report).unwrap_or(Value::Null),
);
}
Ok(None) => {
data.insert("error".into(), Value::from("no shared engine"));
}
Err(e) => {
data.insert("error".into(), Value::from(format!("cycle panicked: {e}")));
}
}
log.append(EventKind::EvolutionTriggered, None, None, data);
}
}))
}
#[cfg(test)]
mod tests {
use super::*;
fn ev(kind: EventKind, action: Option<&str>) -> Event {
Event {
kind,
action_id: action.map(str::to_string),
proposal_id: None,
data: HashMap::new(),
timestamp: chrono::Utc::now(),
prev_hash: None,
hash: None,
}
}
fn connector(slug: &str, connected: bool) -> ConnectorStatus {
ConnectorStatus {
slug: slug.into(),
name: slug.into(),
url: format!("https://example.com/{slug}"),
connected,
tool_count: 1,
enabled_count: 1,
last_error: if connected {
None
} else {
Some("dial failed".into())
},
}
}
#[test]
fn harness_pressure_counts_recurring_failure_share() {
let mut events = vec![
ev(EventKind::ActionRejected, Some("a1")),
ev(EventKind::ActionRejected, Some("a1")),
ev(EventKind::ActionRejected, Some("a1")),
ev(EventKind::ActionRejected, Some("a1")),
];
for _ in 0..4 {
events.push(ev(EventKind::ActionSucceeded, Some("ok")));
}
let c = harness_component_from_events(&events).expect("component");
assert_eq!(c.component, EvolvableComponent::Harness);
assert!((c.signals.pressure - 0.5).abs() < 1e-9, "{c:?}");
assert_eq!(c.signals.evidence, 8);
}
#[test]
fn harness_one_off_failures_are_zero_pressure() {
let events = vec![
ev(EventKind::ActionRejected, Some("a1")),
ev(EventKind::ActionFailed, Some("a2")),
ev(EventKind::ActionSucceeded, Some("a3")),
];
let c = harness_component_from_events(&events).unwrap();
assert_eq!(c.signals.pressure, 0.0);
assert_eq!(c.signals.evidence, 3);
}
#[test]
fn harness_empty_log_is_absent_not_zero() {
assert!(harness_component_from_events(&[]).is_none());
}
#[test]
fn tools_pressure_is_disconnected_share() {
let list = vec![
connector("up", true),
connector("down1", false),
connector("down2", false),
connector("up2", true),
];
let c = tools_component_from_connectors(&list).expect("component");
assert_eq!(c.component, EvolvableComponent::Tools);
assert!((c.signals.pressure - 0.5).abs() < 1e-9, "{c:?}");
assert_eq!(c.signals.evidence, 4);
}
#[test]
fn tools_absent_when_no_connectors_configured() {
assert!(tools_component_from_connectors(&[]).is_none());
}
#[test]
fn failed_trace_events_fold_failure_kinds_only() {
let mut failed = ev(EventKind::ActionFailed, Some("a1"));
failed.data.insert("tool".into(), Value::from("http_get"));
failed.data.insert("error".into(), Value::from("timeout"));
let events = vec![
failed,
ev(EventKind::ActionSucceeded, Some("a2")),
ev(EventKind::PolicyViolation, Some("a3")),
ev(EventKind::ReplanExhausted, None),
];
let traces = failed_trace_events(&events);
assert_eq!(traces.len(), 3);
assert_eq!(traces[0].kind, "action_failed");
assert_eq!(traces[0].tool.as_deref(), Some("http_get"));
assert_eq!(traces[0].action_id.as_deref(), Some("a1"));
assert_eq!(traces[0].reward, Some(0.0));
assert_eq!(traces[1].kind, "policy_violation");
assert_eq!(traces[2].kind, "replan_exhausted");
}
#[test]
fn failed_trace_events_fold_only_ungrounded_completions() {
let mut ungrounded = ev(EventKind::GoalEvaluated, None);
ungrounded.data.insert("met".into(), Value::Bool(true));
ungrounded
.data
.insert("grounded".into(), Value::Bool(false));
let mut grounded = ev(EventKind::GoalEvaluated, None);
grounded.data.insert("met".into(), Value::Bool(true));
grounded.data.insert("grounded".into(), Value::Bool(true));
let mut in_progress = ev(EventKind::GoalEvaluated, None);
in_progress.data.insert("met".into(), Value::Bool(false));
in_progress
.data
.insert("grounded".into(), Value::Bool(false));
let traces = failed_trace_events(&[ungrounded, grounded, in_progress]);
assert_eq!(
traces.len(),
1,
"only the met-but-ungrounded completion is a failure"
);
assert_eq!(traces[0].kind, "goal_evaluated");
assert_eq!(traces[0].reward, Some(0.0));
}
#[test]
fn failed_trace_events_fold_problematic_turn_completions_only() {
let mut truncated = ev(EventKind::TurnCompleted, None);
truncated
.data
.insert("decision".into(), Value::from("empty_tool_calls"));
truncated
.data
.insert("was_truncated".into(), Value::Bool(true));
let mut capped = ev(EventKind::TurnCompleted, None);
capped
.data
.insert("decision".into(), Value::from("max_turns"));
capped
.data
.insert("was_truncated".into(), Value::Bool(false));
let mut clean = ev(EventKind::TurnCompleted, None);
clean
.data
.insert("decision".into(), Value::from("empty_tool_calls"));
clean
.data
.insert("was_truncated".into(), Value::Bool(false));
let traces = failed_trace_events(&[truncated, capped, clean]);
assert_eq!(traces.len(), 2, "a clean finish is not a failure");
assert!(traces.iter().all(|t| t.kind == "turn_completed"));
}
#[test]
fn maintenance_input_prices_backlog_off_live_stats() {
let stats = car_memgine::memsys::MemoryStats {
total_facts: 100,
outstanding_outdated: 5,
facts_superseded: 10,
..Default::default()
};
let input = maintenance_input_from_stats(&stats);
assert_eq!(input.dirty_regions, 15);
assert_eq!(input.total_regions, 100);
assert!((input.global_structural_gain - 0.10).abs() < 1e-9);
let d = decide_maintenance(&input);
assert_eq!(
d.strategy,
car_memgine::maintenance::MaintenanceStrategy::Localized,
"{d:?}"
);
}
#[test]
fn maintenance_input_clean_store_is_noop() {
let input = maintenance_input_from_stats(&car_memgine::memsys::MemoryStats::default());
let d = decide_maintenance(&input);
assert_eq!(
d.strategy,
car_memgine::maintenance::MaintenanceStrategy::NoOp
);
}
#[tokio::test]
async fn memory_evolution_dry_run_reports_without_consolidating() {
let engine = Arc::new(tokio::sync::Mutex::new(MemgineEngine::new(None)));
let out = run_memory_evolution(&engine, true).await.unwrap();
assert!(out.summary.starts_with("dry_run"), "{out:?}");
assert!(!out.applied, "dry run must not count as applied (S2)");
let real = run_memory_evolution(&engine, false).await.unwrap();
assert!(
real.summary.contains("\"mechanism\":\"consolidate\""),
"{real:?}"
);
assert!(real.applied);
}
#[tokio::test]
async fn skills_evolution_without_inference_is_an_honest_error() {
let engine = Arc::new(tokio::sync::Mutex::new(MemgineEngine::new(None)));
let err = run_skills_evolution(&engine, &[], false).await.unwrap_err();
assert_eq!(err, "no inference engine");
}
#[test]
fn context_backoff_widens_exponentially_and_clears_when_the_change_pays() {
let fp = "context:context:abcd1234";
let mut b = ContextBackoff::default();
assert!(b.is_due(fp, 1));
b.note_falsified(fp, 1);
assert!(!b.is_due(fp, 2), "tick 2 still backing off");
assert!(b.is_due(fp, 3));
b.note_falsified(fp, 3);
assert!(!b.is_due(fp, 6));
assert!(b.is_due(fp, 7));
assert!(b.is_due("context:context:99999999", 4));
b.clear(fp);
assert!(b.is_due(fp, 4));
b.note_falsified(fp, 4);
assert!(b.is_due(fp, 6), "exponent restarted at 2^1");
}
#[test]
fn context_backoff_helpers_are_transparent_without_a_backoff() {
assert!(backoff_due(None, "context:context:abcd1234"));
note_falsified(None, "context:context:abcd1234");
clear_backoff(None, "context:context:abcd1234");
}
#[test]
fn skills_backoff_widens_exponentially_and_resets_on_recovery() {
let mut b = SkillsBackoff::default();
let flagged = vec!["web".to_string()];
assert_eq!(b.due(&flagged, 1), flagged);
b.note_attempt("web", 1);
assert!(b.due(&flagged, 2).is_empty(), "tick 2 still backing off");
assert_eq!(b.due(&flagged, 3), flagged);
b.note_attempt("web", 3);
assert!(b.due(&flagged, 6).is_empty());
assert_eq!(b.due(&flagged, 7), flagged);
b.reset_recovered(&[]);
b.note_attempt("web", 10);
assert_eq!(b.due(&flagged, 12), flagged);
}
#[test]
fn skills_backoff_exponent_is_capped() {
let mut b = SkillsBackoff::default();
for t in 0..20 {
b.note_attempt("stuck", t);
}
let flagged = vec!["stuck".to_string()];
assert!(b.due(&flagged, 19 + 63).is_empty());
assert_eq!(b.due(&flagged, 19 + 64), flagged);
}
#[test]
fn skills_backoff_only_gates_the_attempted_domain() {
let mut b = SkillsBackoff::default();
let flagged = vec!["a".to_string(), "b".to_string()];
b.note_attempt("a", 1);
assert_eq!(b.due(&flagged, 2), vec!["b".to_string()]);
}
#[test]
fn cycle_guard_blocks_overlap_and_releases_on_drop() {
let guard = CycleGuard::default();
let token = guard.try_begin().expect("first claim");
assert!(guard.try_begin().is_none(), "in-flight cycle must block");
drop(token);
assert!(guard.try_begin().is_some(), "released after drop");
}
#[test]
fn cycle_guard_releases_even_when_cycle_errors() {
let guard = CycleGuard::default();
let r: Result<(), ()> = {
let _token = guard.try_begin().unwrap();
Err(())
};
assert!(r.is_err());
assert!(guard.try_begin().is_some(), "drop on error path releases");
}
}