From 788ea86e794f1c052a8330c4087247afe49549d2 Mon Sep 17 00:00:00 2001 From: npub1mn7jgtj4w2pd0g0zeuhxsa6jy6p0rewxz4kujt98my82ahfmp72sxjexk7 Date: Mon, 29 Jun 2026 16:49:08 -0400 Subject: [PATCH] feat(acp): add dream consolidation dispatch and preemption Implement Phase 2 of the dream skill: harness-side dispatch logic for memory consolidation turns. Dream is lowest-priority (pending work > heartbeat > dream), preemptible by any inbound event, and follows the heartbeat session lifecycle (reuse until invalidated on cancel). Key additions: - KIND_DREAM_DUE (24300) ephemeral event constant in buzz-core - PromptSource::Dream variant with full match coverage - Dream session management (create once, reuse, invalidate on cancel) - dispatch_dream() with control_tx for preemption - Dream-due event detection sets pending flag (dispatch on next tick) - cancel_in_flight_dream() fires Cancel on any accepted inbound event - Panic recovery clears dream_in_flight via is_dream on TaskMeta - load_dream_prompt() reads .agents/skills/dream/SKILL.md at startup Co-authored-by: Will Pfleger Signed-off-by: Will Pfleger --- crates/buzz-acp/src/lib.rs | 159 +++++++++++++++++++++++++++++++++-- crates/buzz-acp/src/pool.rs | 77 +++++++++++++++-- crates/buzz-core/src/kind.rs | 5 ++ 3 files changed, 228 insertions(+), 13 deletions(-) diff --git a/crates/buzz-acp/src/lib.rs b/crates/buzz-acp/src/lib.rs index 80c9e9d15..e8a21800f 100644 --- a/crates/buzz-acp/src/lib.rs +++ b/crates/buzz-acp/src/lib.rs @@ -16,8 +16,8 @@ use std::time::Duration; use acp::{AcpClient, EnvVar, McpServer}; use anyhow::Result; use buzz_core::kind::{ - KIND_MEMBER_ADDED_NOTIFICATION, KIND_MEMBER_REMOVED_NOTIFICATION, KIND_STREAM_MESSAGE, - KIND_STREAM_REMINDER, KIND_WORKFLOW_APPROVAL_REQUESTED, + KIND_DREAM_DUE, KIND_MEMBER_ADDED_NOTIFICATION, KIND_MEMBER_REMOVED_NOTIFICATION, + KIND_STREAM_MESSAGE, KIND_STREAM_REMINDER, KIND_WORKFLOW_APPROVAL_REQUESTED, }; use buzz_core::observer::{ decrypt_observer_payload, encrypt_observer_payload, OBSERVER_FRAME_TELEMETRY, @@ -1310,6 +1310,7 @@ async fn tokio_main() -> Result<()> { .as_deref() .and_then(|hex| nostr::PublicKey::from_hex(hex).ok()), memory_enabled: config.memory_enabled, + dream_prompt: load_dream_prompt(), }); if !config.memory_enabled { @@ -1329,6 +1330,9 @@ async fn tokio_main() -> Result<()> { None }; let mut heartbeat_in_flight = false; + let mut dream_in_flight = false; + // Whether a dream-due signal has been received and is pending dispatch. + let mut dream_pending = false; let mut presence_heartbeat = if config.presence_enabled { let interval = Duration::from_secs(60); @@ -1687,6 +1691,16 @@ async fn tokio_main() -> Result<()> { continue; } + // 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). + if kind_u32 == KIND_DREAM_DUE { + tracing::info!(target: "dream", "received dream-due signal from relay"); + dream_pending = true; + continue; + } + if config.ignore_self && buzz_event.event.pubkey.to_hex() == pubkey_hex { tracing::debug!(channel_id = %buzz_event.channel_id, "dropping self-authored event"); continue; @@ -1868,6 +1882,12 @@ async fn tokio_main() -> Result<()> { tokio::spawn(async move { pool::reaction_add(&rc, &eid, "👀").await; }); + // Preempt in-flight dream: any real inbound event + // takes priority. Cancel the dream task so the agent + // returns to the pool and can service this event. + if dream_in_flight { + cancel_in_flight_dream(&mut pool); + } } // Event is already queued. If mode requires it AND // the channel has an in-flight task, fire cancel — @@ -1944,8 +1964,10 @@ async fn tokio_main() -> Result<()> { { typing_channels.insert(channel_id, thread_tags); } - } else if pool.any_idle() { + } 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); } else { tracing::debug!("heartbeat_skipped_busy"); } @@ -2013,6 +2035,7 @@ async fn tokio_main() -> Result<()> { &config, *result, &mut heartbeat_in_flight, + &mut dream_in_flight, &removed_channels, &mut crash_history, &respawn_tx, @@ -2027,6 +2050,7 @@ async fn tokio_main() -> Result<()> { &mut queue, &config, &mut heartbeat_in_flight, + &mut dream_in_flight, &removed_channels, &mut typing_channels, &mut crash_history, @@ -2049,6 +2073,7 @@ async fn tokio_main() -> Result<()> { &config, join_error, &mut heartbeat_in_flight, + &mut dream_in_flight, &removed_channels, &mut typing_channels, &mut crash_history, @@ -2359,6 +2384,25 @@ fn signal_in_flight_task( false } +/// Cancel an in-flight dream task so the agent can service real work. +/// +/// Finds the dream task in the pool's task map (identified by `is_dream`) +/// and fires `ControlSignal::Cancel`. The cancel path returns the agent to +/// the pool and clears `dream_in_flight` when the result arrives. +fn cancel_in_flight_dream(pool: &mut AgentPool) { + let entry = pool + .task_map_mut() + .values_mut() + .find(|m| m.is_dream); + + if let Some(meta) = entry { + if let Some(tx) = meta.control_tx.take() { + let _ = tx.send(ControlSignal::Cancel); + tracing::info!(target: "dream", "cancelled in-flight dream (preempted by inbound event)"); + } + } +} + /// Attempt the non-cancelling (ACP) steer for a freshly-queued event. /// /// Caller invariants: @@ -2533,6 +2577,7 @@ fn dispatch_pending( ctx_clone, result_tx, Some(control_rx), + None, ) .await; }); @@ -2545,6 +2590,7 @@ fn dispatch_pending( recoverable_batch, control_tx: Some(control_tx), steer_tx, + is_dream: false, }, ); dispatched_channels.push((channel_id, typing_scope)); @@ -2564,6 +2610,7 @@ fn handle_prompt_result( config: &Config, mut result: PromptResult, heartbeat_in_flight: &mut bool, + dream_in_flight: &mut bool, removed_channels: &HashSet, crash_history: &mut [SlotCircuit], respawn_tx: &mpsc::Sender, @@ -2611,6 +2658,7 @@ fn handle_prompt_result( match &result.source { PromptSource::Channel(ch) => queue.mark_complete(*ch), PromptSource::Heartbeat => *heartbeat_in_flight = false, + PromptSource::Dream => *dream_in_flight = false, } // Strip sessions for channels the agent was removed from while this @@ -2631,7 +2679,7 @@ fn handle_prompt_result( let channel_id = match &result.source { PromptSource::Channel(ch) => Some(*ch), - PromptSource::Heartbeat => None, + PromptSource::Heartbeat | PromptSource::Dream => None, }; let emit_turn_error = |error_msg: &str| { if let Some(ref observer) = observer { @@ -2764,6 +2812,7 @@ fn recover_panicked_agent( config: &Config, join_error: tokio::task::JoinError, heartbeat_in_flight: &mut bool, + dream_in_flight: &mut bool, removed_channels: &HashSet, typing_channels: &mut HashMap, crash_history: &mut [SlotCircuit], @@ -2797,6 +2846,9 @@ fn recover_panicked_agent( queue.mark_complete(ch); typing_channels.remove(&ch); tracing::warn!("cleared wedged in-flight channel {ch} from panicked agent {i}"); + } else if meta.is_dream { + *dream_in_flight = false; + tracing::warn!("cleared wedged dream_in_flight from panicked agent {i}"); } else { *heartbeat_in_flight = false; tracing::warn!("cleared wedged heartbeat_in_flight from panicked agent {i}"); @@ -2859,6 +2911,7 @@ fn drain_ready_join_results( queue: &mut EventQueue, config: &Config, heartbeat_in_flight: &mut bool, + dream_in_flight: &mut bool, removed_channels: &HashSet, typing_channels: &mut HashMap, crash_history: &mut [SlotCircuit], @@ -2875,6 +2928,7 @@ fn drain_ready_join_results( config, join_error, heartbeat_in_flight, + dream_in_flight, removed_channels, typing_channels, crash_history, @@ -2916,7 +2970,7 @@ fn dispatch_heartbeat( let agent_index = agent.index; let abort_handle = pool.join_set.spawn(async move { - pool::run_prompt_task(agent, None, Some(prompt_text), ctx_clone, result_tx, None).await; + pool::run_prompt_task(agent, None, Some(prompt_text), ctx_clone, result_tx, None, None).await; }); pool.task_map_mut().insert( @@ -2927,6 +2981,7 @@ fn dispatch_heartbeat( recoverable_batch: None, control_tx: None, steer_tx: None, + is_dream: false, }, ); *heartbeat_in_flight = true; @@ -2952,6 +3007,96 @@ fn default_heartbeat_prompt() -> String { ) } +/// Load the dream consolidation prompt from the skill file. +/// +/// 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. +fn load_dream_prompt() -> Option { + let path = std::path::Path::new(".agents/skills/dream/SKILL.md"); + match std::fs::read_to_string(path) { + Ok(content) => { + tracing::info!( + target: "dream", + bytes = content.len(), + "loaded dream skill prompt" + ); + Some(content) + } + Err(e) if e.kind() == std::io::ErrorKind::NotFound => { + tracing::debug!(target: "dream", "dream skill not found at {}, dream disabled", path.display()); + None + } + Err(e) => { + tracing::warn!(target: "dream", "failed to read dream skill at {}: {e}", path.display()); + None + } + } +} + +/// Dispatch a dream consolidation turn on an idle agent. +/// +/// Dream is lowest priority — only fires when no pending work AND no heartbeat +/// in flight. Preemptible via `control_tx` (any inbound event cancels it). +fn dispatch_dream( + pool: &mut AgentPool, + ctx: &Arc, + dream_in_flight: &mut bool, + dream_pending: &mut bool, +) { + if *dream_in_flight { + return; + } + let agent = match pool.try_claim(None) { + Some(a) => a, + None => return, + }; + + let prompt_text = match ctx.dream_prompt.as_ref() { + Some(p) => p.clone(), + None => { + // No dream prompt configured — return agent and clear pending. + pool.return_agent(agent); + *dream_pending = false; + return; + } + }; + + let result_tx = pool.result_tx(); + let ctx_clone = Arc::clone(ctx); + let agent_index = agent.index; + + let (control_tx, control_rx) = tokio::sync::oneshot::channel::(); + + let abort_handle = pool.join_set.spawn(async move { + pool::run_prompt_task( + agent, + None, + Some(prompt_text), + ctx_clone, + result_tx, + Some(control_rx), + Some(pool::PromptSource::Dream), + ) + .await; + }); + + pool.task_map_mut().insert( + abort_handle.id(), + pool::TaskMeta { + agent_index, + channel_id: None, + recoverable_batch: None, + control_tx: Some(control_tx), + steer_tx: None, + is_dream: true, + }, + ); + *dream_in_flight = true; + *dream_pending = false; + tracing::info!(agent = agent_index, "dream_fired"); +} + /// Spawn a background respawn task for a crashed agent slot. /// /// Does the circuit breaker check synchronously (non-blocking), then spawns @@ -3370,6 +3515,7 @@ mod owner_control_command_tests { recoverable_batch: None, control_tx: Some(control_tx), steer_tx: None, + is_dream: false, }, ); @@ -3892,12 +4038,14 @@ mod error_outcome_emission_tests { recoverable_batch: None, control_tx: None, steer_tx: None, + is_dream: false, }, ); let mut queue = EventQueue::new(config::DedupMode::Queue); let config = test_config(); let mut heartbeat_in_flight = false; + let mut dream_in_flight = false; let removed_channels = HashSet::new(); let mut crash_history = vec![SlotCircuit { crash_times: Vec::new(), @@ -3921,6 +4069,7 @@ mod error_outcome_emission_tests { &config, result, &mut heartbeat_in_flight, + &mut dream_in_flight, &removed_channels, &mut crash_history, &respawn_tx, diff --git a/crates/buzz-acp/src/pool.rs b/crates/buzz-acp/src/pool.rs index 9a5a872b4..79dd64d24 100644 --- a/crates/buzz-acp/src/pool.rs +++ b/crates/buzz-acp/src/pool.rs @@ -59,6 +59,8 @@ pub struct TaskMeta { /// tasks only — all prompt tasks install a steer channel regardless /// of the agent's name. pub steer_tx: Option>, + /// True when this task is a dream consolidation turn. + pub is_dream: bool, } /// Agent-level model capabilities. Populated on first session creation. @@ -81,11 +83,14 @@ pub struct SessionState { /// channel_id → session_id pub sessions: HashMap, pub heartbeat_session: Option, + pub dream_session: Option, /// Per-channel turn counters for proactive session rotation. /// Incremented on each successful prompt; reset when the session is rotated. pub turn_counts: HashMap, /// Turn counter for the heartbeat session. pub heartbeat_turn_count: u32, + /// Turn counter for the dream session. + pub dream_turn_count: u32, /// channel_id → rendered NIP-AE core prompt section, populated once at /// session creation per Tyler's spec (no mid-session refresh). pub core_sections: HashMap, @@ -102,6 +107,10 @@ impl SessionState { self.heartbeat_session = None; self.heartbeat_turn_count = 0; } + PromptSource::Dream => { + self.dream_session = None; + self.dream_turn_count = 0; + } } } @@ -119,6 +128,8 @@ impl SessionState { self.turn_counts.clear(); self.heartbeat_session = None; self.heartbeat_turn_count = 0; + self.dream_session = None; + self.dream_turn_count = 0; self.core_sections.clear(); } @@ -166,11 +177,12 @@ pub struct PromptResult { pub batch: Option, } -/// Whether the prompt came from a channel event or a heartbeat. +/// Whether the prompt came from a channel event, a heartbeat, or a dream consolidation. #[derive(Debug)] pub enum PromptSource { Channel(Uuid), Heartbeat, + Dream, } /// Apply state effects for Race 1, where a control signal arrives just after the @@ -371,6 +383,9 @@ pub struct PromptContext { /// `[Agent Memory — core]` section. On by default; disabled via /// `--no-memory` / `BUZZ_ACP_NO_MEMORY`. pub memory_enabled: bool, + /// Dream consolidation prompt text, loaded from the skill file at startup. + /// `None` when the skill file doesn't exist — dream dispatch is a no-op. + pub dream_prompt: Option, } impl AgentPool { @@ -905,16 +920,20 @@ pub async fn run_prompt_task( ctx: Arc, result_tx: mpsc::UnboundedSender, control_rx: Option>, + source_hint: Option, ) { - // Is this a channel prompt or a heartbeat? - let source = match &batch { - Some(b) => PromptSource::Channel(b.channel_id), - None => PromptSource::Heartbeat, + // Determine prompt source: explicit hint (dream), or infer from batch. + let source = match source_hint { + Some(s) => s, + None => match &batch { + Some(b) => PromptSource::Channel(b.channel_id), + None => PromptSource::Heartbeat, + }, }; let turn_id = uuid::Uuid::new_v4().to_string(); let observer_channel_id = match &source { PromptSource::Channel(channel_id) => Some(*channel_id), - PromptSource::Heartbeat => None, + PromptSource::Heartbeat | PromptSource::Dream => None, }; agent.acp.set_observer_context(observer::context_for( observer_channel_id, @@ -931,6 +950,7 @@ pub async fn run_prompt_task( "source": match &source { PromptSource::Channel(_) => "channel", PromptSource::Heartbeat => "heartbeat", + PromptSource::Dream => "dream", }, "triggeringEventIds": triggering_event_ids, }), @@ -1019,10 +1039,10 @@ pub async fn run_prompt_task( } // The core section to fold into the system prompt for this turn's session. - // Channel-scoped; heartbeats carry no owner core. + // Channel-scoped; heartbeats and dreams carry no owner core. let agent_core: Option = match &source { PromptSource::Channel(cid) => agent.state.core_sections.get(cid).cloned(), - PromptSource::Heartbeat => None, + PromptSource::Heartbeat | PromptSource::Dream => None, }; let (session_id, is_new_session) = match &source { @@ -1099,6 +1119,42 @@ pub async fn run_prompt_task( } } } + PromptSource::Dream => { + if let Some(sid) = &agent.state.dream_session { + (sid.clone(), false) + } else { + match create_session_and_apply_model(&mut agent, &ctx, None).await { + Ok(sid) => { + tracing::info!( + target: "pool::session", + "created dream session {sid} for agent {}", + agent.index + ); + agent.state.dream_session = Some(sid.clone()); + (sid, true) + } + Err(AcpError::AgentExited) => { + agent.state.invalidate_all(); + let _ = result_tx.send(PromptResult { + agent, + source, + outcome: PromptOutcome::AgentExited, + batch: None, + }); + return; + } + Err(e) => { + let _ = result_tx.send(PromptResult { + agent, + source, + outcome: PromptOutcome::Error(e), + batch: None, + }); + return; + } + } + } + } }; agent.acp.set_observer_context(observer::context_for( observer_channel_id, @@ -1493,6 +1549,10 @@ pub async fn run_prompt_task( agent.state.heartbeat_turn_count += 1; agent.state.heartbeat_turn_count >= limit } + PromptSource::Dream => { + agent.state.dream_turn_count += 1; + agent.state.dream_turn_count >= limit + } } } else { false @@ -2216,6 +2276,7 @@ fn log_stop_reason(source: &PromptSource, stop_reason: &StopReason) { let label = match source { PromptSource::Channel(cid) => format!("channel {cid}"), PromptSource::Heartbeat => "heartbeat".to_string(), + PromptSource::Dream => "dream".to_string(), }; match stop_reason { StopReason::EndTurn => { diff --git a/crates/buzz-core/src/kind.rs b/crates/buzz-core/src/kind.rs index f2e918424..aabb9d07d 100644 --- a/crates/buzz-core/src/kind.rs +++ b/crates/buzz-core/src/kind.rs @@ -275,6 +275,11 @@ pub const KIND_MESH_CONNECT_REQUEST: u32 = 24621; /// so both ends dial near-simultaneously. Tagged `["p", ]`. /// Never stored; seconds expiry. pub const KIND_MESH_CALL_ME_NOW: u32 = 24622; +/// Ephemeral: dream-due signal (relay → harness). The relay emits this when +/// an agent's memory exceeds configured thresholds and the agent has been idle +/// long enough. The harness dispatches a dream consolidation turn in response. +/// Never stored; single-delivery to the authenticated agent. +pub const KIND_DREAM_DUE: u32 = 24300; // Stream messaging /// NIP-29 group chat message kind. V1 used kind:10001 (replaceable range — wrong), then 40001.