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
266 changes: 233 additions & 33 deletions crates/buzz-acp/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -632,8 +632,9 @@ struct QueuedNormalListenerEvent {
/// Visible message that owns this event's lifecycle reactions: the
/// original for an edit, otherwise the event itself.
reaction_target_id: String,
event_for_steer: nostr::Event,
prompt_tag_for_steer: String,
/// The admitted event as a native steer would render it, including an
/// edit's resolved original-message routing.
steer_event: queue::BatchEvent,
}

impl QueuedNormalListenerEvent {
Expand Down Expand Up @@ -662,17 +663,19 @@ impl QueuedNormalListenerEvent {
let Some(signal) = mode_gate_signal(handling, &self.effective_author, owner) else {
return;
};
// Edits take the universal cancel+merge path: a native steer body does
// not yet carry the edit's original-message routing, whereas a
// requeued edit is re-dispatched as its own routed batch.
// A native steer keeps the running turn's `<context>`, so it may only
// carry a message that replies in the same thread. Under the channel
// policy one session spans threads; a message for another thread takes
// the cancel+merge path, whose re-prompt carries its own `<context>`.
let same_reply_thread = queue.in_flight_reply_thread(&self.scope)
== Some(self.steer_event.reply_thread().as_str());
let native_attempted = matches!(signal, ControlSignal::Steer)
&& queue::edit_target_id(&self.event_for_steer).is_none()
&& same_reply_thread
&& try_native_steer(
pool,
queue,
self.scope.clone(),
self.event_for_steer,
self.prompt_tag_for_steer,
self.steer_event,
steer_ack_tx,
);
Comment thread
loganj marked this conversation as resolved.
if !native_attempted {
Expand Down Expand Up @@ -719,14 +722,19 @@ impl NormalListenerIngress {
edit,
} = self;
let reaction_target_id = queue::reaction_target_id(&buzz_event.event);
let event_for_steer = buzz_event.event.clone();
let prompt_tag_for_steer = prompt_tag.clone();
let received_at = std::time::Instant::now();
let steer_event = queue::BatchEvent {
event: buzz_event.event.clone(),
prompt_tag: prompt_tag.clone(),
received_at,
edit: edit.clone(),
};
let channel_id = buzz_event.channel_id;
let accepted = queue.push(QueuedEvent {
channel_id,
scope: session_scope.clone(),
event: buzz_event.event,
received_at: std::time::Instant::now(),
received_at,
prompt_tag,
edit,
});
Expand All @@ -735,8 +743,7 @@ impl NormalListenerIngress {
scope: session_scope,
effective_author,
reaction_target_id,
event_for_steer,
prompt_tag_for_steer,
steer_event,
}
}
}
Expand Down Expand Up @@ -4326,8 +4333,7 @@ fn try_native_steer(
pool: &mut AgentPool,
queue: &mut EventQueue,
scope: scope::SessionScope,
event: nostr::Event,
prompt_tag: String,
be: queue::BatchEvent,
steer_ack_tx: &mpsc::UnboundedSender<SteerAckEvent>,
) -> bool {
let channel_id = scope.channel_id();
Expand All @@ -4344,22 +4350,11 @@ fn try_native_steer(
// channel context and the actor's profile in the original prompt,
// duplicating it here would defeat the point of non-cancelling
// steering (which is to inject only what's new).
let (tag, closing) = queue::native_steer_framing();
let event_id_hex = event.id.to_hex();
let be = queue::BatchEvent {
event,
prompt_tag: prompt_tag.clone(),
received_at: std::time::Instant::now(),
edit: None,
};
let event_block = queue::format_event_block(channel_id, None, &be, None);
let new_message = prompt_framing::semantic_section(tag, "");
let event_section = prompt_framing::semantic_section_with_attributes(
"buzz-event",
&[("type", prompt_tag.as_str())],
&event_block,
);
let body = format!("{new_message}\n\n{event_section}\n\n{closing}");
// The caller steers natively only a message in the running turn's reply
// thread, so the turn's own `<context>` still routes the reply. An edit's
// block names its original (`Edit of:`) and the original's thread root.
let event_id_hex = be.event.id.to_hex();
let body = native_steer_body(channel_id, &be);

let (ack_tx, ack_rx) = tokio::sync::oneshot::channel::<pool::SteerAck>();
let request = pool::SteerRequest {
Expand Down Expand Up @@ -4426,6 +4421,19 @@ fn try_native_steer(
}
}

/// Render the prompt delta sent by [`try_native_steer`].
fn native_steer_body(channel_id: Uuid, be: &queue::BatchEvent) -> String {
let (tag, closing) = queue::native_steer_framing();
let event_block = queue::format_event_block(channel_id, None, be, None);
let new_message = prompt_framing::semantic_section(tag, "");
let event_section = prompt_framing::semantic_section_with_attributes(
"buzz-event",
&[("type", be.prompt_tag.as_str())],
&event_block,
);
format!("{new_message}\n\n{event_section}\n\n{closing}")
}

// ── try_native_steer fallback-log tests ───────────────────────────────────────
//
// Regression for the production tracing event `try_native_steer`'s Err arm
Expand Down Expand Up @@ -4594,8 +4602,12 @@ mod try_native_steer_fallback_log_tests {
pool,
&mut queue,
busy_scope.clone(),
event,
"mention".into(),
queue::BatchEvent {
event,
prompt_tag: "mention".into(),
received_at: std::time::Instant::now(),
edit: None,
},
&steer_ack_tx,
)
});
Expand Down Expand Up @@ -9949,6 +9961,194 @@ mod edit_mention_admission_tests {
}
}

#[cfg(test)]
mod edit_native_steer_tests {
use super::*;
use crate::edit_routing::test_support::{edit_event, message};

/// A steered edit renders its original's routing, not the edit's own
/// bare `e` tag, so the live turn replies in the original's thread.
#[test]
fn native_steer_body_carries_edit_original_routing() {
let root = "ab".repeat(32);
let original = message(Some(&root));
let edit = edit_event(&original.id.to_hex(), &[]);
let be = queue::BatchEvent {
event: edit.clone(),
prompt_tag: "@mention".into(),
received_at: std::time::Instant::now(),
edit: Some(queue::ResolvedEdit {
target_event_id: original.id.to_hex(),
target_thread_tags: queue::parse_thread_tags(&original),
}),
};
let body = native_steer_body(Uuid::new_v4(), &be);
assert!(
body.contains(&format!("Edit of: {}", original.id.to_hex())),
"{body}"
);
assert!(body.contains(&format!("root={root}")), "{body}");
}

/// The listener hands the steer path the same resolved route it queues.
#[test]
fn listener_steer_event_keeps_resolved_edit() {
let original = message(None);
let edit = edit_event(&original.id.to_hex(), &[]);
let channel_id = Uuid::new_v4();
let resolved = queue::ResolvedEdit {
target_event_id: original.id.to_hex(),
target_thread_tags: queue::parse_thread_tags(&original),
};
let ingress = NormalListenerIngress {
buzz_event: relay::BuzzEvent {
connection_generation: 0,
channel_id,
event: edit,
},
effective_author: "author".into(),
prompt_tag: "@mention".into(),
edit: Some(resolved.clone()),
};
let scope = ingress.session_scope(scope::SessionPolicy::Channel, false);
let mut queue = EventQueue::new(config::DedupMode::Queue);
let queued = ingress.push(&mut queue, scope);
assert_eq!(queued.steer_event.edit, Some(resolved));
assert_eq!(queued.reaction_target_id, original.id.to_hex());
}

/// Drive a routed edit of `original` through the listener's steer decision
/// while a turn for `running_event` is in flight under the channel policy.
/// Returns the native steer request, if any, and the control signal sent
/// to the running turn, if any.
fn steer_edit_into_running_turn(
running_event: nostr::Event,
original: &nostr::Event,
) -> (Option<pool::SteerRequest>, Option<ControlSignal>) {
let channel_id = Uuid::new_v4();
let ingress =
|event: nostr::Event, edit: Option<queue::ResolvedEdit>| NormalListenerIngress {
buzz_event: relay::BuzzEvent {
connection_generation: 0,
channel_id,
event,
},
effective_author: "author".into(),
prompt_tag: "@mention".into(),
edit,
};

let mut queue = EventQueue::new(config::DedupMode::Queue);
let running = ingress(running_event, None);
let scope = running.session_scope(scope::SessionPolicy::Channel, false);
running.push(&mut queue, scope.clone());
queue.flush_next().expect("running turn");
assert!(queue.is_scope_in_flight(&scope));

let mut pool = AgentPool::from_slots(vec![]);
let (control_tx, mut control_rx) = tokio::sync::oneshot::channel();
let (steer_tx, mut steer_rx) = tokio::sync::mpsc::channel(1);
let task = pool.join_set.spawn(async {});
pool.task_map_mut().insert(
task.id(),
pool::TaskMeta {
agent_index: 0,
channel_id: Some(channel_id),
scope: Some(scope.clone()),
turn_id: "t".into(),
recoverable_batch: None,
control_tx: Some(control_tx),
steer_tx: Some(steer_tx),
successful_steer_deliveries: HashSet::new(),
},
);

let edit = edit_event(&original.id.to_hex(), &[]);
let resolved = queue::ResolvedEdit {
target_event_id: original.id.to_hex(),
target_thread_tags: queue::parse_thread_tags(original),
};
let edit_ingress = ingress(edit, Some(resolved));
assert_eq!(
edit_ingress.session_scope(scope::SessionPolicy::Channel, false),
scope,
"channel policy: one session spans every thread"
);
let (ack_tx, _ack_rx) = mpsc::unbounded_channel();
edit_ingress
.push(&mut queue, scope.clone())
.steer_or_interrupt(
MultipleEventHandling::Steer,
None,
&mut pool,
&mut queue,
&ack_tx,
);
(steer_rx.try_recv().ok(), control_rx.try_recv().ok())
}

/// A routed edit whose original is in the running turn's thread is
/// steered natively, carrying the original's route, and the running turn
/// is not cancelled.
#[tokio::test]
async fn routed_edit_in_running_thread_steers_natively() {
let root = "ab".repeat(32);
let original = message(Some(&root));
let (steer, control) = steer_edit_into_running_turn(message(Some(&root)), &original);

let request = steer.expect("edit is sent as a native steer");
let body = request.prompt_blocks.join("\n");
assert!(
body.contains(&format!("Edit of: {}", original.id.to_hex())),
"{body}"
);
assert!(body.contains(&format!("root={root}")), "{body}");
assert_eq!(
control, None,
"native steer must not cancel the running turn"
);
}

/// Editing the top-level message a turn is working on steers that turn:
/// replies to the edit belong to the same new thread.
#[tokio::test]
async fn edit_of_running_top_level_trigger_steers_natively() {
let original = message(None);
let (steer, control) = steer_edit_into_running_turn(original.clone(), &original);

assert!(steer.is_some(), "edit is sent as a native steer");
assert_eq!(
control, None,
"native steer must not cancel the running turn"
);
}

/// Under the channel policy a turn started from thread A can be running
/// when an edit whose original is in thread B arrives. A native steer
/// would leave the turn's `<context>` replying to A, so the edit takes the
/// cancel+merge path instead; its re-prompt routes replies to B.
#[tokio::test]
async fn routed_edit_for_another_thread_cancels_and_merges() {
let thread_b = "ab".repeat(32);
let original = message(Some(&thread_b));
let (steer, control) = steer_edit_into_running_turn(message(None), &original);

assert!(steer.is_none(), "no native steer into thread A's turn");
assert_eq!(control, Some(ControlSignal::Steer));
}

/// The same rule covers a top-level original: replies to it open a thread
/// rooted at the original, not at the running turn's top-level trigger.
#[tokio::test]
async fn routed_edit_for_another_top_level_message_cancels_and_merges() {
let original = message(None);
let (steer, control) = steer_edit_into_running_turn(message(None), &original);

assert!(steer.is_none(), "no native steer across top-level threads");
assert_eq!(control, Some(ControlSignal::Steer));
}
}

#[cfg(test)]
mod error_outcome_emission_tests {
//! Pins the policy that error-class outcomes surface to the activity feed
Expand Down
Loading
Loading