use super::*;
pub(super) fn tee_delta_sender(mut senders: Vec<DeltaSender>) -> DeltaSender {
if senders.len() == 1 {
return senders.pop().expect("len checked above");
}
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<String>();
tokio::task::spawn_local(async move {
while let Some(delta) = rx.recv().await {
for sender in &senders {
let _ = sender.send(delta.clone());
}
}
});
tx
}
#[cfg(test)]
pub(super) async fn run_detector_loop_with_sink(
detector_ctx: StreamingDetectorContext,
mut rx: tokio::sync::mpsc::UnboundedReceiver<String>,
mut sink: impl FnMut(&crate::agent_events::AgentEvent),
) {
let mut detector = crate::llm::tools::StreamingToolCallDetector::new(
detector_ctx.session_id,
detector_ctx.known_tools,
);
while let Some(delta) = rx.recv().await {
for event in detector.push(&delta) {
sink(&event);
}
}
for event in detector.finalize() {
sink(&event);
}
}
pub(super) async fn run_detector_loop(
detector_ctx: StreamingDetectorContext,
mut rx: tokio::sync::mpsc::UnboundedReceiver<String>,
mut first_token: super::first_token::FirstTokenTimer,
) {
let mut detector = crate::llm::tools::StreamingToolCallDetector::new(
detector_ctx.session_id,
detector_ctx.known_tools,
);
while let Some(delta) = rx.recv().await {
first_token.observe_delta();
for event in detector.push(&delta) {
crate::agent_events::emit_event(&event);
}
}
for event in detector.finalize() {
crate::agent_events::emit_event(&event);
}
}
pub(crate) const DEFAULT_LLM_CALL_BACKOFF_MS: u64 = 250;
pub(super) const EMPTY_COMPLETION_BUILTIN_RETRIES: usize = 1;