From 9eda16dc7c9b32f438c16610391a3c6e4ff97ebd Mon Sep 17 00:00:00 2001 From: Larry <627498bd4bd1f281a16431e3c6cce3b5c25b6692798c78672298aefbf2f8f8b5@buzz.block.builderlab.xyz> Date: Wed, 30 Sep 2026 12:42:34 -0400 Subject: [PATCH 1/3] fix(acp): steer live turns with routed edits An edit that arrives during an agent turn now uses the native steer path like any other event. The steer carries the edit's resolved original-message route, so the steered turn replies and reacts in the original's thread instead of treating the edit as a new top-level message. The route is resolved at admission, so native steer needs no separate preparation, reservation or membership fence. Fallback, cancellation and channel removal reuse the existing withheld-steer lifecycle, which already keeps the queued edit and its route together. Replaces the earlier 22-commit native-lifecycle slice (old head 82e041561), which resolved edits asynchronously after admission. Signed-off-by: Larry <627498bd4bd1f281a16431e3c6cce3b5c25b6692798c78672298aefbf2f8f8b5@buzz.block.builderlab.xyz> --- crates/buzz-acp/src/lib.rs | 127 +++++++++++++++++++++++++++---------- 1 file changed, 94 insertions(+), 33 deletions(-) diff --git a/crates/buzz-acp/src/lib.rs b/crates/buzz-acp/src/lib.rs index 3f74500d711..0ee62a96734 100644 --- a/crates/buzz-acp/src/lib.rs +++ b/crates/buzz-acp/src/lib.rs @@ -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 { @@ -662,17 +663,12 @@ 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. let native_attempted = matches!(signal, ControlSignal::Steer) - && queue::edit_target_id(&self.event_for_steer).is_none() && try_native_steer( pool, queue, self.scope.clone(), - self.event_for_steer, - self.prompt_tag_for_steer, + self.steer_event, steer_ack_tx, ); if !native_attempted { @@ -719,14 +715,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, }); @@ -735,8 +736,7 @@ impl NormalListenerIngress { scope: session_scope, effective_author, reaction_target_id, - event_for_steer, - prompt_tag_for_steer, + steer_event, } } } @@ -4326,8 +4326,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, ) -> bool { let channel_id = scope.channel_id(); @@ -4344,22 +4343,10 @@ 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}"); + // An edit's block carries its resolved original-message routing, so a + // steered edit is anchored exactly as a dispatched one would be. + 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::(); let request = pool::SteerRequest { @@ -4426,6 +4413,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 @@ -4594,8 +4594,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, ) }); @@ -9949,6 +9953,63 @@ 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()); + } +} + #[cfg(test)] mod error_outcome_emission_tests { //! Pins the policy that error-class outcomes surface to the activity feed From 9b8439473a6157eb7681c151c59813b5c0815c72 Mon Sep 17 00:00:00 2001 From: Larry <627498bd4bd1f281a16431e3c6cce3b5c25b6692798c78672298aefbf2f8f8b5@buzz.block.builderlab.xyz> Date: Wed, 30 Sep 2026 17:00:31 -0400 Subject: [PATCH 2/3] test(acp): drive a routed edit through the native steer decision The new test pushes an admitted edit into a running scope and calls the listener's steer decision. It checks that the sent steer request carries the original's route and that the running turn gets no cancel signal. Restoring the old edit gate, or dropping the route at the send, fails it. Signed-off-by: Larry <627498bd4bd1f281a16431e3c6cce3b5c25b6692798c78672298aefbf2f8f8b5@buzz.block.builderlab.xyz> --- crates/buzz-acp/src/lib.rs | 83 ++++++++++++++++++++++++++++++++++++++ 1 file changed, 83 insertions(+) diff --git a/crates/buzz-acp/src/lib.rs b/crates/buzz-acp/src/lib.rs index 0ee62a96734..a567c6d226e 100644 --- a/crates/buzz-acp/src/lib.rs +++ b/crates/buzz-acp/src/lib.rs @@ -10008,6 +10008,89 @@ mod edit_native_steer_tests { assert_eq!(queued.steer_event.edit, Some(resolved)); assert_eq!(queued.reaction_target_id, original.id.to_hex()); } + + /// End to end through the listener's steer decision: a routed edit that + /// arrives during a running turn goes out as a native steer carrying the + /// original's route, and the running turn is not cancelled. + #[tokio::test] + async fn routed_edit_steers_running_turn_with_original_route() { + let root = "ab".repeat(32); + let original = message(Some(&root)); + let channel_id = Uuid::new_v4(); + let ingress = + |event: nostr::Event, edit: Option| NormalListenerIngress { + buzz_event: relay::BuzzEvent { + connection_generation: 0, + channel_id, + event, + }, + effective_author: "author".into(), + prompt_tag: "@mention".into(), + edit, + }; + + // A turn is already running in the scope. + let mut queue = EventQueue::new(config::DedupMode::Queue); + let running = ingress(message(None), 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 + ); + 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, + ); + + let request = steer_rx.try_recv().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!( + matches!( + control_rx.try_recv(), + Err(tokio::sync::oneshot::error::TryRecvError::Empty) + ), + "native steer must not cancel the running turn" + ); + } } #[cfg(test)] From a41e15c2ad4ecbb75d4f493ca964ddf8fa54ec3f Mon Sep 17 00:00:00 2001 From: Larry <627498bd4bd1f281a16431e3c6cce3b5c25b6692798c78672298aefbf2f8f8b5@buzz.block.builderlab.xyz> Date: Fri, 2 Oct 2026 09:51:24 -0400 Subject: [PATCH 3/3] fix(acp): steer natively only within the running turn's reply thread A native steer adds a message to a running turn without a new , so the turn keeps replying where its own points. Under the channel session policy one session spans threads, so a message for another thread (an edit of a thread-B message during a thread-A turn) would get its reply in thread A. Record each in-flight turn's reply thread in the queue and steer natively only a message in that thread. Any other message takes the cancel+merge path, whose re-prompt carries the message's own full with channel, project, and profile metadata. The reply thread uses the same rule as the thread-session key. Signed-off-by: Larry <627498bd4bd1f281a16431e3c6cce3b5c25b6692798c78672298aefbf2f8f8b5@buzz.block.builderlab.xyz> --- crates/buzz-acp/src/lib.rs | 94 ++++++++++++++++++++++++++++-------- crates/buzz-acp/src/queue.rs | 50 +++++++++++++++++++ crates/buzz-acp/src/scope.rs | 7 +-- 3 files changed, 127 insertions(+), 24 deletions(-) diff --git a/crates/buzz-acp/src/lib.rs b/crates/buzz-acp/src/lib.rs index a567c6d226e..c621785f3b5 100644 --- a/crates/buzz-acp/src/lib.rs +++ b/crates/buzz-acp/src/lib.rs @@ -663,7 +663,14 @@ impl QueuedNormalListenerEvent { let Some(signal) = mode_gate_signal(handling, &self.effective_author, owner) else { return; }; + // A native steer keeps the running turn's ``, 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 ``. + 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) + && same_reply_thread && try_native_steer( pool, queue, @@ -4343,8 +4350,9 @@ 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). - // An edit's block carries its resolved original-message routing, so a - // steered edit is anchored exactly as a dispatched one would be. + // The caller steers natively only a message in the running turn's reply + // thread, so the turn's own `` 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); @@ -10009,13 +10017,14 @@ mod edit_native_steer_tests { assert_eq!(queued.reaction_target_id, original.id.to_hex()); } - /// End to end through the listener's steer decision: a routed edit that - /// arrives during a running turn goes out as a native steer carrying the - /// original's route, and the running turn is not cancelled. - #[tokio::test] - async fn routed_edit_steers_running_turn_with_original_route() { - let root = "ab".repeat(32); - let original = message(Some(&root)); + /// 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, Option) { let channel_id = Uuid::new_v4(); let ingress = |event: nostr::Event, edit: Option| NormalListenerIngress { @@ -10029,9 +10038,8 @@ mod edit_native_steer_tests { edit, }; - // A turn is already running in the scope. let mut queue = EventQueue::new(config::DedupMode::Queue); - let running = ingress(message(None), None); + 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"); @@ -10058,12 +10066,13 @@ mod edit_native_steer_tests { 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), + 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 + scope, + "channel policy: one session spans every thread" ); let (ack_tx, _ack_rx) = mpsc::unbounded_channel(); edit_ingress @@ -10075,22 +10084,69 @@ mod edit_native_steer_tests { &mut queue, &ack_tx, ); + (steer_rx.try_recv().ok(), control_rx.try_recv().ok()) + } - let request = steer_rx.try_recv().expect("edit is sent as a native steer"); + /// 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!( - matches!( - control_rx.try_recv(), - Err(tokio::sync::oneshot::error::TryRecvError::Empty) - ), + 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 `` 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)] diff --git a/crates/buzz-acp/src/queue.rs b/crates/buzz-acp/src/queue.rs index e82a62a04b5..a5b4b9d3f31 100644 --- a/crates/buzz-acp/src/queue.rs +++ b/crates/buzz-acp/src/queue.rs @@ -149,6 +149,11 @@ impl BatchEvent { pub fn routing_thread_tags(&self) -> ThreadTags { routing_thread_tags(&self.event, self.edit.as_ref()) } + + /// See [`reply_thread`]. + pub fn reply_thread(&self) -> String { + reply_thread(&self.event, self.edit.as_ref()) + } } /// Why a batch's prior turn was cancelled — controls how `format_prompt` @@ -236,6 +241,9 @@ pub struct EventQueue { in_flight_deadlines: HashMap, /// Number of events in each in-flight batch (for expiry logging). in_flight_batch_sizes: HashMap, + /// Reply thread of each in-flight turn: the [`reply_thread`] of the batch + /// event whose `` routes the turn's replies (its last event). + in_flight_reply_threads: HashMap, retry_after: HashMap, /// Per-scope retry attempt counter for exponential backoff / dead-lettering. retry_counts: HashMap, @@ -277,6 +285,7 @@ impl EventQueue { in_flight_scopes: HashSet::new(), in_flight_deadlines: HashMap::new(), in_flight_batch_sizes: HashMap::new(), + in_flight_reply_threads: HashMap::new(), retry_after: HashMap::new(), retry_counts: HashMap::new(), dedup_mode, @@ -424,6 +433,7 @@ impl EventQueue { ); self.in_flight_scopes.remove(&scope); self.in_flight_deadlines.remove(&scope); + self.in_flight_reply_threads.remove(&scope); // Recover any withheld goose-native steer events for the expired // scope back to the queue front so normal dispatch delivers // them. Unlike the in-flight batch above (already delivered to a @@ -467,6 +477,7 @@ impl EventQueue { .insert(scope.clone(), now + self.in_flight_deadline); self.in_flight_batch_sizes .insert(scope.clone(), cancelled.len()); + self.record_in_flight_reply_thread(&scope, &cancelled); return Some(FlushBatch { channel_id: scope.channel_id(), scope, @@ -504,6 +515,7 @@ impl EventQueue { .insert(scope.clone(), now + self.in_flight_deadline); self.in_flight_batch_sizes .insert(scope.clone(), events.len()); + self.record_in_flight_reply_thread(&scope, &events); // Merge any cancelled events stored by requeue_as_cancelled(). let cancelled_events = self.cancelled_batches.remove(&scope).unwrap_or_default(); @@ -556,6 +568,7 @@ impl EventQueue { self.in_flight_scopes.remove(&scope); self.in_flight_deadlines.remove(&scope); self.in_flight_batch_sizes.remove(&scope); + self.in_flight_reply_threads.remove(&scope); let now = Instant::now(); match self.retry_after.get(&scope) { // Active throttle → scope was requeued; keep retry_counts intact. @@ -774,6 +787,7 @@ impl EventQueue { ); self.in_flight_scopes.remove(&scope); self.in_flight_deadlines.remove(&scope); + self.in_flight_reply_threads.remove(&scope); // Symmetric with the flush_next expiry block: recover withheld // goose-native steer events for the expired scope so they are // not permanently orphaned in the side table. @@ -929,6 +943,29 @@ impl EventQueue { self.in_flight_scopes.contains(&scope.into_scope()) } + /// The reply thread of the turn in flight for `scope`, if any. + /// + /// A native steer adds a message to a running turn without a new + /// ``, so the turn keeps replying where its own `` + /// points. A message whose [`reply_thread`] differs (possible under the + /// channel session policy) must not be steered natively; the cancel+merge + /// path re-dispatches it with its own full ``. + pub fn in_flight_reply_thread(&self, scope: &SessionScope) -> Option<&str> { + self.in_flight_reply_threads.get(scope).map(String::as_str) + } + + fn record_in_flight_reply_thread(&mut self, scope: &SessionScope, events: &[BatchEvent]) { + match events.last() { + Some(last) => { + self.in_flight_reply_threads + .insert(scope.clone(), last.reply_thread()); + } + None => { + self.in_flight_reply_threads.remove(scope); + } + } + } + /// Whether any scope currently has a turn in flight. pub fn has_in_flight(&self) -> bool { !self.in_flight_scopes.is_empty() @@ -1175,6 +1212,19 @@ pub(crate) fn reaction_target_id(event: &Event) -> String { edit_target_id(event).unwrap_or_else(|| event.id.to_hex()) } +/// The thread that replies to `event` belong to: its routed thread root, or +/// the routed event itself when it is top-level (a reply opens a thread +/// rooted there). Lowercase, so equivalent hex spellings compare equal. +/// +/// This is the thread-session key, and it decides whether a message may be +/// steered natively into a running turn (see [`EventQueue::in_flight_reply_thread`]). +pub(crate) fn reply_thread(event: &Event, edit: Option<&ResolvedEdit>) -> String { + routing_thread_tags(event, edit) + .root_event_id + .unwrap_or_else(|| reaction_target_id(event)) + .to_ascii_lowercase() +} + /// Thread tags that route replies for `event`. See /// [`BatchEvent::routing_thread_tags`]. pub(crate) fn routing_thread_tags(event: &Event, edit: Option<&ResolvedEdit>) -> ThreadTags { diff --git a/crates/buzz-acp/src/scope.rs b/crates/buzz-acp/src/scope.rs index 1ec23ebeb7e..a0dd2176fb0 100644 --- a/crates/buzz-acp/src/scope.rs +++ b/crates/buzz-acp/src/scope.rs @@ -23,7 +23,7 @@ use nostr::Event; use uuid::Uuid; -use crate::queue::{reaction_target_id, routing_thread_tags, ResolvedEdit}; +use crate::queue::{reply_thread, ResolvedEdit}; /// Operator policy controlling how ACP provider sessions are scoped. /// @@ -150,12 +150,9 @@ impl SessionScope { return Self::Conversation { channel_id }; } - let root_event_id = routing_thread_tags(event, edit) - .root_event_id - .unwrap_or_else(|| reaction_target_id(event)); Self::Thread { channel_id, - root_event_id: root_event_id.to_ascii_lowercase(), + root_event_id: reply_thread(event, edit), } }