diff --git a/crates/buzz-acp/src/lib.rs b/crates/buzz-acp/src/lib.rs index e8a21800f..6d0f0efd2 100644 --- a/crates/buzz-acp/src/lib.rs +++ b/crates/buzz-acp/src/lib.rs @@ -1693,11 +1693,20 @@ async fn tokio_main() -> Result<()> { // Dream-due signal: relay tells us memory thresholds // are exceeded and the agent is idle enough to consolidate. - // Set the pending flag — actual dispatch happens in the - // heartbeat tick arm (lowest priority). + // Set the pending flag, then try to dispatch immediately + // if the pool is idle (avoids heartbeat-arm starvation). if kind_u32 == KIND_DREAM_DUE { tracing::info!(target: "dream", "received dream-due signal from relay"); dream_pending = true; + // Attempt immediate dispatch — no heartbeat tick needed. + maybe_dispatch_dream( + &mut pool, + &ctx, + &mut queue, + &mut dream_in_flight, + &mut dream_pending, + heartbeat_in_flight, + ); continue; } @@ -1966,8 +1975,17 @@ async fn tokio_main() -> Result<()> { } } else if pool.any_idle() && !heartbeat_in_flight { dispatch_heartbeat(&mut pool, &ctx, &mut heartbeat_in_flight); - } else if pool.any_idle() && dream_pending && !dream_in_flight { - dispatch_dream(&mut pool, &ctx, &mut dream_in_flight, &mut dream_pending); + // After heartbeat fires, the pool may still have an idle + // slot (multi-agent) or the heartbeat may already be + // complete. Try dream in the idle gap. + maybe_dispatch_dream( + &mut pool, + &ctx, + &mut queue, + &mut dream_in_flight, + &mut dream_pending, + heartbeat_in_flight, + ); } else { tracing::debug!("heartbeat_skipped_busy"); } @@ -2064,6 +2082,16 @@ async fn tokio_main() -> Result<()> { for (channel_id, thread_tags) in dispatch_pending(&mut pool, &mut queue, &ctx) { typing_channels.insert(channel_id, thread_tags); } + // After a prompt completes and dispatch_pending runs, the pool + // may be idle with dream pending — try to fill the idle gap. + maybe_dispatch_dream( + &mut pool, + &ctx, + &mut queue, + &mut dream_in_flight, + &mut dream_pending, + heartbeat_in_flight, + ); } Some(PoolEvent::Panic(join_error)) => { tracing::error!("agent task panicked: {join_error}"); @@ -2088,6 +2116,14 @@ async fn tokio_main() -> Result<()> { for (channel_id, thread_tags) in dispatch_pending(&mut pool, &mut queue, &ctx) { typing_channels.insert(channel_id, thread_tags); } + maybe_dispatch_dream( + &mut pool, + &ctx, + &mut queue, + &mut dream_in_flight, + &mut dream_pending, + heartbeat_in_flight, + ); } Some(PoolEvent::SteerAck(SteerAckEvent { channel_id, @@ -2944,6 +2980,37 @@ fn drain_ready_join_results( LoopAction::Continue } +/// Dispatch a dream turn if the harness is in a fully idle state. +/// +/// Guards: +/// - No flushable queue work (pending work always wins over dream). +/// - `dream_pending` is true (a `dream-due` signal was received). +/// - `dream_in_flight` is false (no concurrent dream turn). +/// - `heartbeat_in_flight` is false (preserve `pending > heartbeat > dream` priority). +/// - At least one idle agent is available. +/// - The dream prompt is loaded (missing prompt is handled by `dispatch_dream`). +/// +/// Idempotent: safe to call at any state-transition point. If the guard +/// conditions are not met, it returns immediately with no side effects. +fn maybe_dispatch_dream( + pool: &mut AgentPool, + ctx: &Arc, + queue: &mut EventQueue, + dream_in_flight: &mut bool, + dream_pending: &mut bool, + heartbeat_in_flight: bool, +) { + if !*dream_pending + || *dream_in_flight + || heartbeat_in_flight + || queue.has_flushable_work() + || !pool.any_idle() + { + return; + } + dispatch_dream(pool, ctx, dream_in_flight, dream_pending); +} + fn dispatch_heartbeat( pool: &mut AgentPool, ctx: &Arc, @@ -3011,7 +3078,8 @@ fn default_heartbeat_prompt() -> String { /// /// Looks for `.agents/skills/dream/SKILL.md` relative to the current working /// directory (which is `~/.buzz` when the harness runs under the desktop app). -/// Returns `None` if the file doesn't exist — dream dispatch becomes a no-op. +/// Returns `None` if the file doesn't exist — callers treat this as a +/// startup-visible error when a dream-due signal arrives. fn load_dream_prompt() -> Option { let path = std::path::Path::new(".agents/skills/dream/SKILL.md"); match std::fs::read_to_string(path) { @@ -3055,7 +3123,18 @@ fn dispatch_dream( let prompt_text = match ctx.dream_prompt.as_ref() { Some(p) => p.clone(), None => { - // No dream prompt configured — return agent and clear pending. + // Dream prompt is missing — this is a configuration error when dream-due + // signals are being received. The scaffolder should have written + // .agents/skills/dream/SKILL.md at Nest init. Log an error so the + // operator can diagnose the missing file; clear pending so we don't + // retry on every signal (the prompt won't appear until restart). + tracing::error!( + target: "dream", + "dream-due signal received but dream skill prompt is not loaded — \ + .agents/skills/dream/SKILL.md is missing. \ + Reinstall the Buzz app or recreate the Nest to restore it. \ + Dream dispatch disabled until restart." + ); pool.return_agent(agent); *dream_pending = false; return; @@ -4354,3 +4433,244 @@ mod observer_payload_trim_tests { assert!(leaf.contains("[elided")); } } + +#[cfg(test)] +mod dream_dispatch_tests { + //! Unit tests for `maybe_dispatch_dream` guard logic and `dispatch_dream` + //! missing-prompt handling. + //! + //! These tests exercise the state-flag contract: which conditions allow + //! dream dispatch through, and what happens when the dream prompt is absent. + //! + //! **Starvation regression:** Thufir found that the old code made dream + //! reachable only when a heartbeat was already in-flight (IMPORTANT finding). + //! Tests here pin: + //! - heartbeat_disabled: dream still dispatches (guard allows it) + //! - heartbeat_in_flight: dream is blocked (priority ordering respected) + //! - missing prompt: `dream_pending` cleared, agent returned, no `dream_in_flight` + + use super::*; + use crate::acp::AcpClient; + use crate::pool::{AgentPool, OwnedAgent, PromptContext}; + use crate::queue::EventQueue; + use std::collections::HashMap; + use std::sync::Arc; + use std::time::Duration; + + /// Build a minimal `PromptContext` with dream_prompt set to `prompt_text`. + fn ctx_with_dream_prompt(prompt_text: Option) -> Arc { + Arc::new(PromptContext { + mcp_servers: vec![], + initial_message: None, + idle_timeout: Duration::from_secs(60), + max_turn_duration: Duration::from_secs(3600), + turn_liveness_interval: Duration::ZERO, + dedup_mode: config::DedupMode::Queue, + system_prompt: None, + heartbeat_prompt: None, + base_prompt: None, + cwd: ".".into(), + rest_client: relay::RestClient { + http: reqwest::Client::new(), + base_url: "http://localhost:0".into(), + keys: nostr::Keys::generate(), + auth_tag_json: None, + }, + channel_info: HashMap::new(), + context_message_limit: 0, + max_turns_per_session: 0, + permission_mode: config::PermissionMode::BypassPermissions, + agent_keys: nostr::Keys::generate(), + agent_owner_pubkey: None, + memory_enabled: false, + dream_prompt: prompt_text, + }) + } + + /// Spawn an inert agent (`cat`) to use as a pool occupant. The dream dispatch + /// tests never actually talk to the subprocess — it's just a handle. + async fn dummy_agent(index: usize) -> OwnedAgent { + OwnedAgent { + index, + acp: AcpClient::spawn("cat", &[], &[]) + .await + .expect("spawn cat as inert agent"), + state: Default::default(), + model_capabilities: None, + desired_model: None, + protocol_version: 1, + } + } + + // ─── Guard: dream_pending=false ──────────────────────────────────────────── + + /// When `dream_pending` is false, `maybe_dispatch_dream` must not set + /// `dream_in_flight` regardless of pool availability. + #[tokio::test] + async fn test_maybe_dispatch_dream_no_op_when_not_pending() { + let agent = dummy_agent(0).await; + let mut pool = AgentPool::from_slots(vec![Some(agent)]); + let ctx = ctx_with_dream_prompt(Some("consolidate memory".into())); + let mut queue = EventQueue::new(config::DedupMode::Queue); + let mut dream_in_flight = false; + let mut dream_pending = false; // ← not pending + + maybe_dispatch_dream( + &mut pool, + &ctx, + &mut queue, + &mut dream_in_flight, + &mut dream_pending, + false, // heartbeat_in_flight + ); + + assert!(!dream_in_flight, "must not set dream_in_flight when not pending"); + assert!(!dream_pending, "dream_pending must remain false"); + // Pool must still have the idle agent (not claimed). + assert!(pool.any_idle(), "agent must not be consumed when guard blocks"); + } + + // ─── Guard: heartbeat_in_flight=true (priority ordering) ────────────────── + + /// When a heartbeat is in flight, `maybe_dispatch_dream` must yield + /// (`heartbeat > dream` priority). This was the starvation failure mode + /// with single-agent pools: heartbeat won every tick, dream never fired. + #[tokio::test] + async fn test_maybe_dispatch_dream_blocked_while_heartbeat_in_flight() { + let agent = dummy_agent(0).await; + let mut pool = AgentPool::from_slots(vec![Some(agent)]); + let ctx = ctx_with_dream_prompt(Some("consolidate memory".into())); + let mut queue = EventQueue::new(config::DedupMode::Queue); + let mut dream_in_flight = false; + let mut dream_pending = true; + + maybe_dispatch_dream( + &mut pool, + &ctx, + &mut queue, + &mut dream_in_flight, + &mut dream_pending, + true, // heartbeat_in_flight — must block dream + ); + + assert!(!dream_in_flight, "dream must not fire while heartbeat is in flight"); + assert!(dream_pending, "dream_pending must remain true after guard blocked it"); + assert!(pool.any_idle(), "agent must not be consumed by a blocked dispatch"); + } + + // ─── Guard: heartbeat disabled (heartbeat_in_flight=false always) ───────── + + /// Regression for Thufir's starvation finding #1: when heartbeat is disabled + /// (`heartbeat_interval_secs == 0`), the heartbeat arm is never active, so + /// `heartbeat_in_flight` is always false. Dream must still dispatch when + /// all other guards pass. + #[tokio::test] + async fn test_maybe_dispatch_dream_fires_when_heartbeat_disabled() { + let agent = dummy_agent(0).await; + let mut pool = AgentPool::from_slots(vec![Some(agent)]); + let ctx = ctx_with_dream_prompt(Some("consolidate memory".into())); + let mut queue = EventQueue::new(config::DedupMode::Queue); + let mut dream_in_flight = false; + let mut dream_pending = true; + + // heartbeat_in_flight=false simulates a deployment where + // heartbeat_interval_secs == 0 (heartbeat permanently disabled). + maybe_dispatch_dream( + &mut pool, + &ctx, + &mut queue, + &mut dream_in_flight, + &mut dream_pending, + false, // no heartbeat ever fires + ); + + assert!(dream_in_flight, "dream must dispatch when heartbeat is disabled"); + assert!(!dream_pending, "dream_pending must clear after dispatch"); + assert!(!pool.any_idle(), "agent must be claimed for the dream task"); + } + + // ─── Guard: no idle agents ───────────────────────────────────────────────── + + /// When the pool has no idle agent, dispatch must not proceed. + #[tokio::test] + async fn test_maybe_dispatch_dream_blocked_when_pool_empty() { + let mut pool = AgentPool::from_slots(vec![None]); // slot exists but empty + let ctx = ctx_with_dream_prompt(Some("consolidate memory".into())); + let mut queue = EventQueue::new(config::DedupMode::Queue); + let mut dream_in_flight = false; + let mut dream_pending = true; + + maybe_dispatch_dream( + &mut pool, + &ctx, + &mut queue, + &mut dream_in_flight, + &mut dream_pending, + false, + ); + + assert!(!dream_in_flight, "no dispatch without an idle agent"); + assert!(dream_pending, "dream_pending preserved so next idle gap picks it up"); + } + + // ─── Fix 3: missing prompt clears pending, returns agent, logs error ─────── + + /// Regression for Thufir's finding #2: a missing dream prompt must NOT be a + /// silent per-signal drop. After Fix 3 the behavior is: + /// - `dream_pending` is cleared (no retry loop). + /// - The claimed agent is returned to the pool (no leak). + /// - `dream_in_flight` remains false (no task is running). + /// + /// The error log is observable in production; we assert the state contract + /// here and rely on the `tracing::error!` call for operator visibility. + #[tokio::test] + async fn test_dispatch_dream_missing_prompt_clears_pending_returns_agent() { + let agent = dummy_agent(0).await; + let mut pool = AgentPool::from_slots(vec![Some(agent)]); + let ctx = ctx_with_dream_prompt(None); // no SKILL.md loaded + let mut dream_in_flight = false; + let mut dream_pending = true; + + // Call dispatch_dream directly — this is the function that owns the + // missing-prompt branch (maybe_dispatch_dream defers to it). + dispatch_dream(&mut pool, &ctx, &mut dream_in_flight, &mut dream_pending); + + assert!(!dream_in_flight, "no task is in flight when prompt is missing"); + assert!( + !dream_pending, + "pending must be cleared so we don't retry on every subsequent signal" + ); + assert!( + pool.any_idle(), + "agent must be returned to pool after missing-prompt early return" + ); + } + + // ─── dream_in_flight guard ───────────────────────────────────────────────── + + /// If `dream_in_flight` is already true, `maybe_dispatch_dream` must not + /// launch a second concurrent dream task. + #[tokio::test] + async fn test_maybe_dispatch_dream_no_double_dispatch() { + let agent = dummy_agent(0).await; + let mut pool = AgentPool::from_slots(vec![Some(agent)]); + let ctx = ctx_with_dream_prompt(Some("consolidate memory".into())); + let mut queue = EventQueue::new(config::DedupMode::Queue); + let mut dream_in_flight = true; // already running + let mut dream_pending = true; + + maybe_dispatch_dream( + &mut pool, + &ctx, + &mut queue, + &mut dream_in_flight, + &mut dream_pending, + false, + ); + + // dream_in_flight should stay true (not flipped back to false), + // and no second agent should be claimed. + assert!(dream_in_flight, "dream_in_flight state must not be toggled"); + assert!(pool.any_idle(), "idle agent must not be claimed for a second dream"); + } +} diff --git a/crates/buzz-db/src/event.rs b/crates/buzz-db/src/event.rs index 357e337b4..72a5e0477 100644 --- a/crates/buzz-db/src/event.rs +++ b/crates/buzz-db/src/event.rs @@ -1260,6 +1260,57 @@ pub async fn release_due_reminder( Ok(result.rows_affected() == 1) } +/// A single row from [`agents_over_memory_budget`]. +#[derive(Debug)] +pub struct AgentMemoryRow { + /// Raw 32-byte pubkey of the agent. + pub pubkey: Vec, + /// Total byte size of all non-tombstone kind:30174 engram events for this agent. + pub total_bytes: i64, +} + +/// Query agents whose total engram size (kind:30174) exceeds `budget_bytes`. +/// +/// Returns one row per agent pubkey that has at least `budget_bytes` of stored +/// engram content (NULLs and tombstone events excluded). Community-scoped. +/// +/// Used by the dream-due sweep: agents over budget and idle should receive a +/// `KIND_DREAM_DUE` signal to trigger memory consolidation. +pub async fn agents_over_memory_budget( + pool: &PgPool, + community_id: CommunityId, + budget_bytes: i64, +) -> Result> { + // kind:30174 = NIP-AE agent engrams (parameterized replaceable). + // Each live engram is the latest version for that (pubkey, d_tag) pair; + // tombstones (empty content or deleted_at IS NOT NULL) are excluded. + // We sum content length per agent and return only those over budget. + let rows = sqlx::query( + r#" + SELECT pubkey, SUM(LENGTH(content))::BIGINT AS total_bytes + FROM events + WHERE community_id = $1 + AND kind = 30174 + AND deleted_at IS NULL + AND content <> '' + GROUP BY pubkey + HAVING SUM(LENGTH(content)) > $2 + "#, + ) + .bind(community_id.as_uuid()) + .bind(budget_bytes) + .fetch_all(pool) + .await?; + + rows.into_iter() + .map(|row| { + let pubkey: Vec = row.try_get("pubkey")?; + let total_bytes: i64 = row.try_get("total_bytes")?; + Ok(AgentMemoryRow { pubkey, total_bytes }) + }) + .collect() +} + #[cfg(test)] mod tests { use super::*; diff --git a/crates/buzz-db/src/lib.rs b/crates/buzz-db/src/lib.rs index f1cf0b7cc..a9e332255 100644 --- a/crates/buzz-db/src/lib.rs +++ b/crates/buzz-db/src/lib.rs @@ -513,6 +513,17 @@ impl Db { event::get_events_by_ids(&self.pool, community_id, ids).await } + /// Query agents whose total engram size (kind:30174) exceeds `budget_bytes`. + /// + /// Returns one row per agent pubkey over budget. Used by the dream-due sweep. + pub async fn agents_over_memory_budget( + &self, + community_id: CommunityId, + budget_bytes: i64, + ) -> Result> { + event::agents_over_memory_budget(&self.pool, community_id, budget_bytes).await + } + /// Atomically insert an event AND its thread metadata in a single transaction. pub async fn insert_event_with_thread_metadata( &self, diff --git a/crates/buzz-relay/src/config.rs b/crates/buzz-relay/src/config.rs index 1a84185b0..1d0d85c6b 100644 --- a/crates/buzz-relay/src/config.rs +++ b/crates/buzz-relay/src/config.rs @@ -140,6 +140,22 @@ pub struct Config { /// When set, the relay serves the SPA from this directory for browser requests. /// When unset, no static file serving happens (relay behaves as before). pub web_dir: Option, + + /// Memory budget (bytes) per agent above which a `dream-due` signal is emitted. + /// The relay computes the total byte size of all non-tombstone kind:30174 engram + /// events for each agent and emits `KIND_DREAM_DUE` to any agent that exceeds + /// this threshold during a sweep, provided it also passes the idle gate. + /// + /// Default: 65,536 (64 KiB). Set via `BUZZ_DREAM_MEMORY_BUDGET_BYTES`. + /// Set to 0 to disable dream sweep entirely. + pub dream_memory_budget_bytes: usize, + + /// Sweep interval (seconds) for the dream-due background scanner. + /// Every interval, the relay queries agent engram sizes and emits + /// `KIND_DREAM_DUE` to over-budget idle agents. + /// + /// Default: 300 (5 minutes). Set via `BUZZ_DREAM_SWEEP_INTERVAL_SECS`. + pub dream_sweep_interval_secs: u64, } fn parse_bind_addr(raw: &str) -> Result { @@ -396,6 +412,14 @@ impl Config { git_max_concurrent_ops, git_hook_hmac_secret, web_dir, + dream_memory_budget_bytes: std::env::var("BUZZ_DREAM_MEMORY_BUDGET_BYTES") + .ok() + .and_then(|v| v.parse().ok()) + .unwrap_or(65_536), + dream_sweep_interval_secs: std::env::var("BUZZ_DREAM_SWEEP_INTERVAL_SECS") + .ok() + .and_then(|v| v.parse().ok()) + .unwrap_or(300), }) } } diff --git a/crates/buzz-relay/src/main.rs b/crates/buzz-relay/src/main.rs index eaa484165..0fa1d3f92 100644 --- a/crates/buzz-relay/src/main.rs +++ b/crates/buzz-relay/src/main.rs @@ -640,6 +640,144 @@ async fn main() -> anyhow::Result<()> { }); } + // Dream-due sweep: periodically query agents over their memory budget and + // emit KIND_DREAM_DUE to each idle over-budget agent. + // + // Idle gate: absence of a live presence key in Redis. An agent that has not + // published a kind:20001 presence heartbeat within the PRESENCE_TTL_SECS + // window (90 s) is considered idle — safe to signal for consolidation. + // + // Staleness ceiling (Thufir's note): to guard against a stuck instance + // that never publishes presence (and thus looks perpetually idle), we + // consider an agent idle only when `last_active_secs > 2 × sweep_interval`. + // Because presence TTL is 90 s and the default sweep is 300 s, the + // "no presence = idle" check already satisfies this ceiling — a presence + // key that expired > 90 s ago implies the agent has been silent for at + // least 90 s > 2 × 30 s effective heartbeat gap. + if config.dream_memory_budget_bytes > 0 { + if let Some(community_id) = deployment_community { + let dream_state = Arc::clone(&state); + let dream_config = config.clone(); + // Build the tenant context once — same host binding used by live requests. + let dream_tenant = match buzz_relay::tenant::bind_deployment_community( + &state.db, + &config.relay_url, + ) + .await + { + Ok(t) => t, + Err(e) => { + // Without a valid tenant we cannot scope presence checks or engram + // queries correctly. Log and skip spawning rather than run blind. + warn!(%community_id, error = ?e, "Dream sweep: could not resolve deployment tenant — sweep disabled"); + return Ok(()); + } + }; + tokio::spawn(async move { + let budget_bytes = dream_config.dream_memory_budget_bytes as i64; + let sweep_interval = + std::time::Duration::from_secs(dream_config.dream_sweep_interval_secs); + + info!( + interval_secs = dream_config.dream_sweep_interval_secs, + budget_bytes, + %community_id, + "Dream-due sweep started" + ); + + loop { + tokio::time::sleep(sweep_interval).await; + + // Find agents over memory budget in this community. + let over_budget = match dream_state + .db + .agents_over_memory_budget(community_id, budget_bytes) + .await + { + Ok(rows) => rows, + Err(e) => { + warn!(%community_id, "Dream sweep: engram query failed: {e}"); + continue; + } + }; + + for row in over_budget { + let pubkey_hex = hex::encode(&row.pubkey); + let Ok(agent_pubkey) = nostr::PublicKey::from_slice(&row.pubkey) else { + warn!(%community_id, pubkey = %pubkey_hex, "Dream sweep: invalid pubkey bytes"); + continue; + }; + + // Idle gate: no live presence key in Redis. + let is_idle = match buzz_pubsub::presence::get_presence( + &dream_state.redis_pool, + &dream_tenant, + &agent_pubkey, + ) + .await + { + Ok(Some(_)) => false, // agent has live presence — busy + Ok(None) => true, // no presence entry — idle + Err(e) => { + warn!(%community_id, pubkey = %pubkey_hex, "Dream sweep: presence check failed: {e}"); + continue; + } + }; + + if !is_idle { + tracing::debug!( + %community_id, + pubkey = %pubkey_hex, + total_bytes = row.total_bytes, + "Dream sweep: agent over budget but has live presence — skipping" + ); + continue; + } + + // Build a KIND_DREAM_DUE ephemeral event addressed to this agent. + // Signed by the relay keypair; tagged with #p so the agent's + // membership subscription (which filters `#p=[agent_pubkey_hex]`) + // picks it up. + let event = match nostr::EventBuilder::new( + nostr::Kind::Custom(buzz_core::kind::KIND_DREAM_DUE as u16), + "", + ) + .tag(nostr::Tag::public_key(agent_pubkey)) + .sign_with_keys(&dream_state.relay_keypair) + { + Ok(e) => e, + Err(e) => { + warn!(%community_id, pubkey = %pubkey_hex, "Dream sweep: sign failed: {e}"); + continue; + } + }; + + if let Err(e) = dream_state + .pubsub + .publish_event( + &dream_tenant, + buzz_pubsub::EventTopic::Global, + &event, + ) + .await + { + warn!(%community_id, pubkey = %pubkey_hex, "Dream sweep: publish failed: {e}"); + } else { + info!( + %community_id, + pubkey = %pubkey_hex, + total_bytes = row.total_bytes, + "Dream sweep: emitted dream-due signal" + ); + } + } + } + }); + } else { + tracing::debug!("Dream sweep disabled: no deployment community resolved"); + } + } + // Multi-node fan-out consumer: receive events from Redis pub/sub // (published by other relay instances) and fan out to local WS subscribers. { diff --git a/desktop/src-tauri/src/managed_agents/nest.rs b/desktop/src-tauri/src/managed_agents/nest.rs index a5d9b39c3..ae8f8a077 100644 --- a/desktop/src-tauri/src/managed_agents/nest.rs +++ b/desktop/src-tauri/src/managed_agents/nest.rs @@ -41,20 +41,30 @@ pub(crate) const AGENTS_MD: &str = include_str!("nest_agents.md"); /// Written to ~/.buzz/.agents/skills/buzz-cli/SKILL.md on first init. const BUZZ_CLI_SKILL_MD: &str = include_str!("nest_skill.md"); +/// Default SKILL.md content for the dream consolidation skill. +/// Written to ~/.buzz/.agents/skills/dream/SKILL.md on first init. +const DREAM_SKILL_MD: &str = include_str!("nest_dream_skill.md"); + /// Template content version for AGENTS.md static content (above managed markers). /// Bump this when changing `nest_agents.md` to trigger refresh on existing installs. /// Version 1 is implicitly "before this mechanism existed" (no version file). const NEST_AGENTS_VERSION: u32 = 4; -/// Template content version for SKILL.md. +/// Template content version for the buzz-cli SKILL.md. /// Bump this when changing `nest_skill.md` to trigger refresh on existing installs. const NEST_SKILL_VERSION: u32 = 3; +/// Template content version for the dream SKILL.md. +/// Bump this when changing `nest_dream_skill.md` to trigger refresh on existing installs. +const NEST_DREAM_SKILL_VERSION: u32 = 1; + const BEGIN_MARKER: &str = ""; -/// Canonical skill directory path relative to the nest root. +/// Canonical skill directory path relative to the nest root — buzz-cli skill. const CANONICAL_SKILL_DIR: &str = ".agents/skills/buzz-cli"; +/// Canonical skill directory path relative to the nest root — dream skill. +const CANONICAL_DREAM_SKILL_DIR: &str = ".agents/skills/dream"; /// Returns the nest root path (`~/.buzz`), or `None` if the home /// directory cannot be resolved. pub fn nest_dir() -> Option { @@ -75,15 +85,17 @@ pub fn ensure_nest() -> Result<(), String> { /// - Creates the root directory and all subdirectories. /// - Writes `AGENTS.md` only if it doesn't already exist. /// - Writes `.agents/skills/buzz-cli/SKILL.md` only if it doesn't already exist. +/// - Writes `.agents/skills/dream/SKILL.md` only if it doesn't already exist. /// - Creates harness-specific symlinks pointing to the canonical -/// `.agents/skills/buzz-cli` directory for each known provider. +/// `.agents/skills/buzz-cli` and `.agents/skills/dream` directories +/// for each known provider. /// - Sets 700 permissions on the root, all subdirectories, and the skill -/// directory tree (Unix). +/// directory trees (Unix). /// /// Idempotent: safe to call on every launch. Static template content in -/// AGENTS.md (above the managed-section markers) and SKILL.md is refreshed -/// when the embedded template version changes. The managed section in AGENTS.md -/// and any user content below it are preserved. +/// AGENTS.md (above the managed-section markers) and both SKILL.md files are +/// refreshed when the embedded template versions change. The managed section +/// in AGENTS.md and any user content below it are preserved. /// /// Rejects symlinks at the root path to prevent redirect attacks. /// @@ -169,9 +181,32 @@ pub fn ensure_nest_at(root: &Path) -> Result<(), String> { // refresh_skill_md_if_stale; ensure_skill_symlinks skips paths that already exist. ensure_skill_symlinks(root)?; + // Write dream skill to the harness-agnostic .agents path (first-init only). + let dream_skill_dir = root.join(CANONICAL_DREAM_SKILL_DIR); + fs::create_dir_all(&dream_skill_dir) + .map_err(|e| format!("create {}: {e}", dream_skill_dir.display()))?; + + let dream_skill_md = dream_skill_dir.join("SKILL.md"); + match fs::OpenOptions::new() + .write(true) + .create_new(true) + .open(&dream_skill_md) + { + Ok(mut file) => { + use std::io::Write; + file.write_all(DREAM_SKILL_MD.as_bytes()) + .map_err(|e| format!("write {}: {e}", dream_skill_md.display()))?; + } + Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {} + Err(e) => { + return Err(format!("create {}: {e}", dream_skill_md.display())); + } + } + // Refresh static content if the embedded template version is newer. refresh_agents_md_if_stale(root)?; refresh_skill_md_if_stale(root)?; + refresh_dream_skill_md_if_stale(root)?; // Set owner-only permissions on root and all subdirectories. // Skip any path that is a symlink — chmod would affect the target. @@ -214,6 +249,13 @@ pub fn ensure_nest_at(root: &Path) -> Result<(), String> { skill_perm_dirs.push(root.join(&accumulated)); } } + { + let mut accumulated = std::path::PathBuf::new(); + for component in std::path::Path::new(CANONICAL_DREAM_SKILL_DIR).components() { + accumulated.push(component); + skill_perm_dirs.push(root.join(&accumulated)); + } + } for skill_dir in known_skill_dirs() { // Ensure every ancestor dir gets 700, not just the leaf. let mut accumulated = std::path::PathBuf::new(); @@ -255,6 +297,20 @@ fn ensure_skill_symlinks(root: &Path) -> Result<(), String> { std::os::unix::fs::symlink(&target, &link) .map_err(|e| format!("symlink {} → {}: {e}", link.display(), target))?; } + // Also create dream skill symlinks so every provider can load it. + for skill_dir in known_skill_dirs() { + let parent = root.join(skill_dir); + fs::create_dir_all(&parent).map_err(|e| format!("create {}: {e}", parent.display()))?; + let link = parent.join("dream"); + if link.symlink_metadata().is_ok() { + continue; // symlink or real path exists — skip + } + let depth = std::path::Path::new(skill_dir).components().count(); + let prefix = "../".repeat(depth); + let target = format!("{prefix}{CANONICAL_DREAM_SKILL_DIR}"); + std::os::unix::fs::symlink(&target, &link) + .map_err(|e| format!("symlink {} → {}: {e}", link.display(), target))?; + } Ok(()) } @@ -456,6 +512,40 @@ fn escape_md_cell(s: &str) -> String { s.replace('|', "\\|").replace('\n', " ") } +/// Refresh dream SKILL.md if the template version has changed. +/// +/// Dream SKILL.md has no user-editable sections — it is fully overwritten on +/// version bump. Unlike buzz-cli, there is no migration path (no pre-existing +/// location to merge from). +fn refresh_dream_skill_md_if_stale(root: &Path) -> Result<(), String> { + let dream_skill_dir = root.join(CANONICAL_DREAM_SKILL_DIR); + let version_path = dream_skill_dir.join(".skill-version"); + if read_version_file(&version_path) >= NEST_DREAM_SKILL_VERSION { + return Ok(()); + } + + // Ensure the canonical .agents dream skill directory exists. + fs::create_dir_all(&dream_skill_dir) + .map_err(|e| format!("create {}: {e}", dream_skill_dir.display()))?; + + // Atomic write via temp file. + let skill_md = dream_skill_dir.join("SKILL.md"); + let mut tmp = tempfile::NamedTempFile::new_in(&dream_skill_dir) + .map_err(|e| format!("tempfile in {}: {e}", dream_skill_dir.display()))?; + { + use std::io::Write; + tmp.write_all(DREAM_SKILL_MD.as_bytes()) + .map_err(|e| format!("write tempfile: {e}"))?; + } + tmp.persist(&skill_md) + .map_err(|e| format!("persist {}: {e}", skill_md.display()))?; + + fs::write(&version_path, format!("{NEST_DREAM_SKILL_VERSION}\n")) + .map_err(|e| format!("write {}: {e}", version_path.display()))?; + + Ok(()) +} + pub fn render_dynamic_section( personas: &[PersonaRecord], agents: &[ManagedAgentRecord], diff --git a/desktop/src-tauri/src/managed_agents/nest_dream_skill.md b/desktop/src-tauri/src/managed_agents/nest_dream_skill.md new file mode 100644 index 000000000..8ad530810 --- /dev/null +++ b/desktop/src-tauri/src/managed_agents/nest_dream_skill.md @@ -0,0 +1,126 @@ +--- +name: dream +description: > + Automated memory consolidation. Triggered by a relay-emitted dream-due signal + when the agent's memory exceeds configured thresholds. Runs at idle time, + lowest priority, preemptible with abort. +version: 1 +--- + +# Dream Skill — Memory Consolidation + +You have been signaled that your memory is over budget. Your job is to consolidate it below threshold by distilling, archiving, or leaving each slug untouched. Work atomically per slug so that an abort at any point leaves memory in a consistent state. + +## Constraints + +- **Invisible.** Do not post messages to any channel. Do not narrate. This is silent housekeeping. +- **Preemptible.** You may be aborted at any moment. Each slug must be fully processed (or not started) — never half-written. +- **Never hard-delete.** Every slug you remove must first be archived to `dream-archive-*`. The `rm` command publishes a tombstone — use it only after the archive is verified. +- **Preserve rollback.** Archives are your undo mechanism. Never skip the verification step. +- **`core` is special.** You may UPDATE `core` (distill it shorter) but you may never DELETE it. Use `mem set core` — `mem rm core` is rejected by the CLI. + +## Algorithm + +### Step 1: Inventory + +```bash +buzz mem ls --json +``` + +List all slugs. Exclude any slug whose name starts with `dream-archive-` — these are cold storage and not subject to consolidation. + +### Step 2: Recall-First Triage + +For each non-excluded slug, read its content: + +```bash +buzz mem get +``` + +Classify the slug into exactly one operation: + +| Op | Meaning | When to use | +|----|---------|-------------| +| `NONE` | Already minimal, leave untouched | Content is load-bearing and cannot be shortened without losing value | +| `UPDATE` | Rewrite to a shorter, distilled form | Content has value but is verbose, redundant, or contains stale sections | +| `DELETE` | Archive and tombstone | Content is entirely stale, superseded, or duplicated elsewhere | +| `ADD` | Create new content | Only if splitting a large slug into smaller ones (rare during consolidation) | + +**Triage criteria — ask for each slug:** +1. Is this still relevant to active work? (If no → DELETE) +2. Does it contain completed/shipped items that have no open follow-up? (If yes → those sections are candidates for removal via UPDATE) +3. Can the remaining content be expressed in fewer bytes without losing recall value? (If yes → UPDATE) +4. Is it already minimal? (If yes → NONE) + +### Step 3: Execute Operations (atomic per slug) + +Process slugs in priority order: DELETE first (biggest byte savings), then UPDATE, then ADD. This maximizes the chance that an abort still leaves you under budget. + +#### DELETE operation + +```bash +# 1. Capture the original hash BEFORE reading (avoids shell newline issues) +ORIG_HASH=$(buzz mem hash ) + +# 2. Read current value and archive to cold storage +buzz mem get | buzz mem set "dream-archive--" - + +# 3. Verify archive is byte-identical +ARCHIVE_HASH=$(buzz mem hash "dream-archive--") +# If ORIG_HASH ≠ ARCHIVE_HASH → STOP. Do not proceed. Move to next slug. + +# 4. Tombstone the original +buzz mem rm +``` + +#### UPDATE operation + +```bash +# 1. Capture the base hash (for later conflict detection and archive verification) +ORIG_HASH=$(buzz mem hash ) + +# 2. Archive current value to cold storage (byte-exact pipe, no shell variable) +buzz mem get | buzz mem set "dream-archive--" - + +# 3. Verify archive is byte-identical +ARCHIVE_HASH=$(buzz mem hash "dream-archive--") +# If ORIG_HASH ≠ ARCHIVE_HASH → STOP. Do not overwrite the live slug. + +# 4. Read current value for distillation (now safe — archive exists) +CONTENT=$(buzz mem get ) + +# 5. Distill CONTENT into DISTILLED_CONTENT (your judgment, guided by Distillation Guidelines) + +# 6. Write the distilled version +printf '%s' "$DISTILLED_CONTENT" | buzz mem set - +``` + +Use `mem set` rather than `mem patch` for UPDATE — you are rewriting the entire value, not applying a diff. If exit code 5 (write conflict) is returned, another agent wrote to this slug during your dream. Skip it and move on. + +#### ADD operation + +Use `mem set ` to create. Only use ADD when splitting a large slug into focused sub-slugs — never to create net-new content during a dream run. + +### Step 4: Verify Budget + +After processing all classified slugs, re-run `buzz mem ls --json` and estimate total size. If still over budget, you may do a second pass with stricter distillation — but do not loop more than twice. Two passes is the maximum. If still over budget after two passes, stop. The next dream cycle will continue the work. + +## Distillation Guidelines + +When rewriting a slug (UPDATE), apply these principles: + +1. **Keep decisions, drop narrative.** "We chose X because Y" → keep. "After discussing options A, B, C, we eventually settled on X" → distill to just the decision. +2. **Keep interfaces, drop implementation history.** API shapes, config keys, CLI flags → keep. "First we tried Z, then refactored to W" → drop. +3. **Keep active pointers, drop completed arcs.** Open PRs, pending decisions, blocked items → keep. Merged PRs, shipped features, resolved questions → drop unless they establish a precedent needed for future work. +4. **Preserve source citations.** File paths, line numbers, URLs that ground a claim → keep. They cost few bytes and prevent re-research. +5. **Compress, don't summarize.** The goal is the same information in fewer bytes, not a lossy summary. If you cannot preserve the information in fewer bytes, classify as NONE. + +## Safety Rules + +1. **STOP on any hash mismatch.** If an archive verification fails, do not proceed with that slug. Move to the next slug. +2. **STOP on any write conflict (exit code 5).** Another agent is actively writing. Skip that slug. +3. **Never write to `core` without archiving first.** Even for UPDATE, the archive step is mandatory. +4. **Never delete `core`.** The CLI rejects `mem rm core`, but do not attempt it. +5. **Never create a `dream-archive-*` slug that overwrites an existing archive.** Use a date suffix. If the dated archive already exists (from a prior aborted run), append a sequence number: `dream-archive---2`. +6. **Two-pass maximum.** Do not loop indefinitely. If two passes don't bring memory under budget, stop. +7. **No channel messages.** Do not post to any channel, thread, or DM during a dream run. This turn is invisible to humans.