Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 14 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -174,6 +174,20 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/).

### Changed

- **DRY**: consolidated the four near-identical memory-maintenance-loop spawn blocks in
`src/runner.rs` (CLI/TUI, previously inline), `src/acp.rs` (`build_acp_deps`, previously
inline), `src/daemon.rs` (`run_daemon`, previously inline), and `src/serve/deps.rs`
(`spawn_memory_maintenance_loops`) into a single shared
`agent_setup::spawn_memory_maintenance_loops`, now called by all four entry points (#6180).
The function takes `status_tx: Option<&UnboundedSender<String>>` to unify the ACP/`/sessions*`
(`None`) vs. CLI/TUI/daemon (`Some`) difference in the hebbian-consolidation loop's status
sender, and a `skip_eviction: bool` preserving `runner.rs`'s pre-existing `--bare` gate, which
only ever applied to the eviction loop. The `acp.rs`/`daemon.rs` unit tests previously
reconstructed a hand-written copy of the production spawn block (unlike `serve/deps.rs`'s
test, which already called the real function); both now call the shared production function
directly via `AppBuilder::for_test`, closing a test-realism gap left by #6170. No behavior
change for any entry point.

- **DRY**: extracted the `zeph worktree clean` reconcile → remove → prune pipeline and its
removed/skipped/errored counting into a single `WorktreeManager::clean` method plus a
`format_clean_summary` free function (`zeph-worktree`), now shared by both the CLI
Expand Down
333 changes: 31 additions & 302 deletions src/acp.rs
Original file line number Diff line number Diff line change
Expand Up @@ -642,16 +642,23 @@ async fn build_acp_deps(

let config = app.config();

// #5914/#5979: memory maintenance loops — mirrors src/runner.rs's CLI/TUI wiring so ACP
// sessions (standalone `--acp` and the ACP half of `serve-sessions --acp`) get the same
// ongoing eviction/tier-promotion/scene-consolidation/consolidation/forgetting/guidelines/
// tree-consolidation/hebbian-consolidation/episodic-consolidation/optical-forgetting sweeps
// instead of an ever-growing, never-maintained memory store. Spawned once per connection
// (shared across all sessions on `acp_mem_supervisor`), matching runner.rs's
// once-per-process cadence. Extracted into `spawn_acp_memory_maintenance_loops` (#6170) so
// the regression test below can call the real, shipped production wiring instead of
// reconstructing a copy of it.
spawn_acp_memory_maintenance_loops(app, &memory, &provider, &acp_mem_supervisor);
// #5914/#5979/#6180: memory maintenance loops, via the shared
// `agent_setup::spawn_memory_maintenance_loops` (also used by `src/runner.rs`,
// `src/daemon.rs`, `src/serve/deps.rs`) so ACP sessions (standalone `--acp` and the ACP half
// of `serve-sessions --acp`) get the same ongoing eviction/tier-promotion/
// scene-consolidation/consolidation/forgetting/guidelines/tree-consolidation/
// hebbian-consolidation/episodic-consolidation/optical-forgetting sweeps instead of an
// ever-growing, never-maintained memory store. Spawned once per connection (shared across
// all sessions on `acp_mem_supervisor`), matching runner.rs's once-per-process cadence.
agent_setup::spawn_memory_maintenance_loops(
app,
&memory,
&provider,
&acp_mem_supervisor,
None,
false,
"acp",
);

let filter_registry = if config.tools.filters.enabled {
zeph_tools::OutputFilterRegistry::default_filters(&config.tools.filters)
Expand Down Expand Up @@ -1187,291 +1194,6 @@ async fn build_acp_deps(
Ok((deps, keepalive))
}

/// Spawns the ten ongoing memory-maintenance sweeps (eviction, tier-promotion,
/// scene-consolidation, consolidation, forgetting — #5914; plus guidelines, tree-consolidation,
/// hebbian-consolidation, episodic-consolidation, optical-forgetting — #5979) on `supervisor`,
/// shared by every session built from this ACP connection's `SharedCore` — mirrors
/// `src/runner.rs`'s CLI/TUI wiring and `src/serve/deps.rs`'s `spawn_memory_maintenance_loops`.
///
/// Extracted from [`build_acp_deps`] (#6170) so the regression test
/// `acp_memory_maintenance_loops_registered_on_connection_supervisor` can call this real,
/// shipped production function directly rather than reconstructing a copy of its spawn blocks —
/// a hand-reconstructed copy could never catch a broken or inverted `if config.memory.*.enabled`
/// guard in production, since it would only prove the test's own copy of the condition was
/// correct.
#[cfg(feature = "acp")]
#[allow(clippy::too_many_lines)]
fn spawn_acp_memory_maintenance_loops(
app: &AppBuilder,
memory: &std::sync::Arc<zeph_memory::semantic::SemanticMemory>,
provider: &zeph_llm::any::AnyProvider,
supervisor: &zeph_common::TaskSupervisor,
) {
let config = app.config();
{
let store = std::sync::Arc::new(memory.sqlite().clone());
let embedding = memory.embedding_store().cloned();
let eviction_cfg = config.memory.eviction.clone();
let policy = std::sync::Arc::new(zeph_memory::EbbinghausPolicy::default());
let cancel = supervisor.cancellation_token();
supervisor.spawn(zeph_common::task_supervisor::TaskDescriptor {
name: "mem-eviction",
restart: zeph_common::task_supervisor::RestartPolicy::RunOnce,
factory: move || {
zeph_memory::start_eviction_loop(
store.clone(),
embedding.clone(),
eviction_cfg.clone(),
policy.clone(),
cancel.clone(),
)
},
});
}
{
let store = std::sync::Arc::new(memory.sqlite().clone());
let tier_cfg = zeph_memory::TierPromotionConfig {
enabled: config.memory.tiers.enabled,
promotion_min_sessions: config.memory.tiers.promotion_min_sessions,
similarity_threshold: config.memory.tiers.similarity_threshold,
sweep_interval_secs: config.memory.tiers.sweep_interval_secs,
sweep_batch_size: config.memory.tiers.sweep_batch_size,
embed_timeout_secs: config.memory.semantic.embed_timeout_secs,
};
let tier_provider = provider.clone();
let cancel = supervisor.cancellation_token();
supervisor.spawn(zeph_common::task_supervisor::TaskDescriptor {
name: "mem-tier-promotion",
restart: zeph_common::task_supervisor::RestartPolicy::RunOnce,
factory: move || {
zeph_memory::start_tier_promotion_loop(
store.clone(),
tier_provider.clone(),
tier_cfg.clone(),
cancel.clone(),
)
},
});
}
{
let store = std::sync::Arc::new(memory.sqlite().clone());
let scene_provider = app
.build_scene_provider()
.unwrap_or_else(|| provider.clone());
let scene_cfg = zeph_memory::SceneConfig {
enabled: config.memory.tiers.scene_enabled,
similarity_threshold: config.memory.tiers.scene_similarity_threshold,
batch_size: config.memory.tiers.scene_batch_size,
sweep_interval_secs: config.memory.tiers.scene_sweep_interval_secs,
};
let cancel = supervisor.cancellation_token();
supervisor.spawn(zeph_common::task_supervisor::TaskDescriptor {
name: "mem-scene-consolidation",
restart: zeph_common::task_supervisor::RestartPolicy::RunOnce,
factory: move || {
zeph_memory::start_scene_consolidation_loop(
store.clone(),
scene_provider.clone(),
scene_cfg.clone(),
cancel.clone(),
)
},
});
}
{
let store = std::sync::Arc::new(memory.sqlite().clone());
let consolidation_cfg = zeph_memory::ConsolidationConfig {
enabled: config.memory.consolidation.enabled,
confidence_threshold: config.memory.consolidation.confidence_threshold,
sweep_interval_secs: config.memory.consolidation.sweep_interval_secs,
sweep_batch_size: config.memory.consolidation.sweep_batch_size,
similarity_threshold: config.memory.consolidation.similarity_threshold,
llm_timeout_secs: config.memory.consolidation.llm_timeout_secs,
embed_timeout_secs: config.memory.semantic.embed_timeout_secs,
};
let consolidation_provider = app
.build_consolidation_provider()
.unwrap_or_else(|| provider.clone());
let cancel = supervisor.cancellation_token();
supervisor.spawn(zeph_common::task_supervisor::TaskDescriptor {
name: "mem-consolidation",
restart: zeph_common::task_supervisor::RestartPolicy::RunOnce,
factory: move || {
zeph_memory::start_consolidation_loop(
store.clone(),
consolidation_provider.clone(),
consolidation_cfg.clone(),
cancel.clone(),
)
},
});
}
{
let store = std::sync::Arc::new(memory.sqlite().clone());
let forgetting_cfg = zeph_memory::ForgettingConfig {
enabled: config.memory.forgetting.enabled,
decay_rate: config.memory.forgetting.decay_rate,
forgetting_floor: config.memory.forgetting.forgetting_floor,
sweep_interval_secs: config.memory.forgetting.sweep_interval_secs,
sweep_batch_size: config.memory.forgetting.sweep_batch_size,
replay_window_hours: config.memory.forgetting.replay_window_hours,
replay_min_access_count: config.memory.forgetting.replay_min_access_count,
protect_recent_hours: config.memory.forgetting.protect_recent_hours,
protect_min_access_count: config.memory.forgetting.protect_min_access_count,
};
let cancel = supervisor.cancellation_token();
supervisor.spawn(zeph_common::task_supervisor::TaskDescriptor {
name: "mem-forgetting",
restart: zeph_common::task_supervisor::RestartPolicy::RunOnce,
factory: move || {
zeph_memory::start_forgetting_loop(
store.clone(),
forgetting_cfg.clone(),
cancel.clone(),
)
},
});
}
if config.memory.compression_guidelines.enabled {
let store = std::sync::Arc::new(memory.sqlite().clone());
let guidelines_provider = app
.build_guidelines_provider()
.unwrap_or_else(|| provider.clone());
let token_counter = std::sync::Arc::clone(&memory.token_counter);
let guidelines_cfg = config.memory.compression_guidelines.clone();
let cancel = supervisor.cancellation_token();
supervisor.spawn(zeph_common::task_supervisor::TaskDescriptor {
name: "mem-guidelines",
restart: zeph_common::task_supervisor::RestartPolicy::RunOnce,
factory: move || {
zeph_memory::start_guidelines_updater(
store.clone(),
guidelines_provider.clone(),
token_counter.clone(),
guidelines_cfg.clone(),
cancel.clone(),
)
},
});
}
if config.memory.tree.enabled {
let store = std::sync::Arc::new(memory.sqlite().clone());
let tree_provider = app
.build_tree_consolidation_provider()
.unwrap_or_else(|| provider.clone());
let tree_cfg = zeph_memory::TreeConsolidationConfig {
enabled: config.memory.tree.enabled,
sweep_interval_secs: config.memory.tree.sweep_interval_secs,
batch_size: config.memory.tree.batch_size,
similarity_threshold: config.memory.tree.similarity_threshold,
max_level: config.memory.tree.max_level,
min_cluster_size: config.memory.tree.min_cluster_size,
embed_timeout_secs: config.memory.semantic.embed_timeout_secs,
};
let cancel = supervisor.cancellation_token();
supervisor.spawn(zeph_common::task_supervisor::TaskDescriptor {
name: "mem-tree-consolidation",
restart: zeph_common::task_supervisor::RestartPolicy::RunOnce,
factory: move || {
zeph_memory::start_tree_consolidation_loop(
store.clone(),
tree_provider.clone(),
tree_cfg.clone(),
cancel.clone(),
)
},
});
}
if config.memory.hebbian.enabled && config.memory.hebbian.consolidation_interval_secs > 0 {
let store = std::sync::Arc::new(memory.sqlite().clone());
let hebbian_consolidation_cfg = zeph_memory::HebbianConsolidationConfig {
consolidation_interval_secs: config.memory.hebbian.consolidation_interval_secs,
consolidation_threshold: config.memory.hebbian.consolidation_threshold,
max_candidates_per_sweep: config.memory.hebbian.max_candidates_per_sweep,
consolidation_cooldown_secs: config.memory.hebbian.consolidation_cooldown_secs,
consolidation_prompt_timeout_secs: config
.memory
.hebbian
.consolidation_prompt_timeout_secs,
consolidation_max_neighbors: config.memory.hebbian.consolidation_max_neighbors,
};
let hebbian_provider = app
.build_hebbian_consolidation_provider()
.unwrap_or_else(|| provider.clone());
let cancel = supervisor.cancellation_token();
supervisor.spawn(zeph_common::task_supervisor::TaskDescriptor {
name: "mem-hebbian-consolidation",
restart: zeph_common::task_supervisor::RestartPolicy::RunOnce,
factory: move || {
zeph_memory::spawn_hebbian_consolidation_loop(
store.clone(),
hebbian_consolidation_cfg.clone(),
hebbian_provider.clone(),
None,
cancel.clone(),
)
},
});
}
if config.memory.episodic_consolidation.enabled {
let store = std::sync::Arc::new(memory.sqlite().clone());
let ep_cfg = zeph_memory::EpisodicConsolidationConfig {
enabled: config.memory.episodic_consolidation.enabled,
consolidation_provider: config
.memory
.episodic_consolidation
.consolidation_provider
.clone(),
interval_secs: config.memory.episodic_consolidation.interval_secs,
batch_size: config.memory.episodic_consolidation.batch_size,
min_age_secs: config.memory.episodic_consolidation.min_age_secs,
dedup_jaccard_threshold: config.memory.episodic_consolidation.dedup_jaccard_threshold,
};
let ep_provider = app
.build_episodic_consolidation_provider()
.unwrap_or_else(|| provider.clone());
let ep_qdrant = memory.embedding_store().cloned();
let cancel = supervisor.cancellation_token();
supervisor.spawn(zeph_common::task_supervisor::TaskDescriptor {
name: "mem-episodic-consolidation",
restart: zeph_common::task_supervisor::RestartPolicy::RunOnce,
factory: move || {
zeph_memory::start_episodic_consolidation_loop(
store.clone(),
ep_provider.clone(),
ep_cfg.clone(),
ep_qdrant.clone(),
cancel.clone(),
)
},
});
}
if config.memory.optical_forgetting.enabled {
let store = std::sync::Arc::new(memory.sqlite().clone());
let optical_provider = app
.build_optical_forgetting_provider()
.unwrap_or_else(|| provider.clone());
let optical_cfg = config.memory.optical_forgetting.clone();
let forgetting_floor = config.memory.forgetting.forgetting_floor;
let cancel = supervisor.cancellation_token();
tracing::info_span!("acp.memory.optical_forgetting.startup").in_scope(|| {
supervisor.spawn(zeph_common::task_supervisor::TaskDescriptor {
name: "mem-optical-forgetting",
restart: zeph_common::task_supervisor::RestartPolicy::RunOnce,
factory: move || {
zeph_memory::start_optical_forgetting_loop(
store.clone(),
optical_provider.clone(),
optical_cfg.clone(),
forgetting_floor,
cancel.clone(),
)
},
});
});
}
}

/// Text shown to the client when session persistence is disabled due to a held write lock.
///
/// Deliberately omits the lock path: it is an absolute filesystem path (leaks the server's
Expand Down Expand Up @@ -3974,12 +3696,11 @@ mod tests {
/// Regression test confirming all ten memory-maintenance loops (eviction, tier-promotion,
/// scene-consolidation, consolidation, forgetting — #5914; plus guidelines,
/// tree-consolidation, hebbian-consolidation, episodic-consolidation, optical-forgetting —
/// #5979) are actually registered on the ACP connection's own `TaskSupervisor`. Calls the
/// real, shipped `spawn_acp_memory_maintenance_loops` (extracted from `build_acp_deps` in
/// #6170) directly rather than reconstructing a copy of its spawn blocks, so a broken or
/// inverted `if config.memory.*.enabled` guard in production is actually caught here instead
/// of only proving the test's own copy of the condition was correct. The five #5979 loops
/// are config-gated, so the config below explicitly enables each.
/// #5979) are actually registered on the ACP connection's own `TaskSupervisor` by the shared
/// `agent_setup::spawn_memory_maintenance_loops` (also called by `build_acp_deps` in
/// production, and by `src/runner.rs`/`src/daemon.rs`/`src/serve/deps.rs`, #6180) — asserts
/// every expected task name is present in the connection supervisor's snapshot. The five
/// #5979 loops are config-gated, so the config below explicitly enables each.
#[tokio::test]
async fn acp_memory_maintenance_loops_registered_on_connection_supervisor() {
let mock_provider =
Expand All @@ -4005,7 +3726,15 @@ mod tests {
let cancel = tokio_util::sync::CancellationToken::new();
let supervisor = zeph_common::TaskSupervisor::new(cancel);

spawn_acp_memory_maintenance_loops(&app, &memory, &mock_provider, &supervisor);
agent_setup::spawn_memory_maintenance_loops(
&app,
&memory,
&mock_provider,
&supervisor,
None,
false,
"acp",
);

let names: std::collections::HashSet<String> = supervisor
.snapshot()
Expand Down
Loading
Loading