gemini_adk_fluent_rs/live/
callbacks.rs1use std::future::Future;
4use std::sync::Arc;
5use std::time::Duration;
6
7use bytes::Bytes;
8
9use gemini_adk_rs::live::CallbackMode;
10use gemini_adk_rs::State;
11use gemini_genai_rs::prelude::*;
12
13use super::Live;
14
15impl Live {
16 pub fn before_tool_response<F, Fut>(mut self, f: F) -> Self
35 where
36 F: Fn(Vec<FunctionResponse>, gemini_adk_rs::State) -> Fut + Send + Sync + 'static,
37 Fut: Future<Output = Vec<FunctionResponse>> + Send + 'static,
38 {
39 self.callbacks.before_tool_response = Some(Arc::new(move |responses, state| {
40 Box::pin(f(responses, state))
41 }));
42 self
43 }
44
45 pub fn on_turn_boundary<F, Fut>(mut self, f: F) -> Self
62 where
63 F: Fn(gemini_adk_rs::State, Arc<dyn gemini_genai_rs::session::SessionWriter>) -> Fut
64 + Send
65 + Sync
66 + 'static,
67 Fut: Future<Output = ()> + Send + 'static,
68 {
69 self.callbacks.on_turn_boundary =
70 Some(Arc::new(move |state, writer| Box::pin(f(state, writer))));
71 self
72 }
73
74 pub fn on_audio(mut self, f: impl Fn(&Bytes) + Send + Sync + 'static) -> Self {
78 self.callbacks.on_audio = Some(Box::new(f));
79 self
80 }
81
82 pub fn on_text(mut self, f: impl Fn(&str) + Send + Sync + 'static) -> Self {
84 self.callbacks.on_text = Some(Box::new(f));
85 self
86 }
87
88 pub fn on_text_complete(mut self, f: impl Fn(&str) + Send + Sync + 'static) -> Self {
90 self.callbacks.on_text_complete = Some(Box::new(f));
91 self
92 }
93
94 pub fn on_input_transcript(mut self, f: impl Fn(&str, bool) + Send + Sync + 'static) -> Self {
96 self.callbacks.on_input_transcript = Some(Box::new(f));
97 self
98 }
99
100 pub fn on_output_transcript(mut self, f: impl Fn(&str, bool) + Send + Sync + 'static) -> Self {
102 self.callbacks.on_output_transcript = Some(Box::new(f));
103 self
104 }
105
106 pub fn on_thought(mut self, f: impl Fn(&str) + Send + Sync + 'static) -> Self {
111 self.callbacks.on_thought = Some(Box::new(f));
112 self
113 }
114
115 pub fn on_vad_start(mut self, f: impl Fn() + Send + Sync + 'static) -> Self {
117 self.callbacks.on_vad_start = Some(Box::new(f));
118 self
119 }
120
121 pub fn on_vad_end(mut self, f: impl Fn() + Send + Sync + 'static) -> Self {
123 self.callbacks.on_vad_end = Some(Box::new(f));
124 self
125 }
126
127 pub fn on_usage(mut self, f: impl Fn(&UsageMetadata) + Send + Sync + 'static) -> Self {
133 self.callbacks.on_usage = Some(Box::new(f));
134 self
135 }
136
137 pub fn on_phase(mut self, f: impl Fn(SessionPhase) + Send + Sync + 'static) -> Self {
142 self.callbacks.on_phase = Some(Box::new(f));
143 self
144 }
145
146 pub fn on_interrupted<F, Fut>(mut self, f: F) -> Self
150 where
151 F: Fn() -> Fut + Send + Sync + 'static,
152 Fut: Future<Output = ()> + Send + 'static,
153 {
154 self.callbacks.on_interrupted = Some(Arc::new(move || Box::pin(f())));
155 self
156 }
157
158 pub fn on_tool_call<F, Fut>(mut self, f: F) -> Self
162 where
163 F: Fn(Vec<FunctionCall>, State) -> Fut + Send + Sync + 'static,
164 Fut: Future<Output = Option<Vec<FunctionResponse>>> + Send + 'static,
165 {
166 self.callbacks.on_tool_call = Some(Arc::new(move |calls, state| Box::pin(f(calls, state))));
167 self
168 }
169
170 pub fn on_tool_cancelled<F, Fut>(mut self, f: F) -> Self
175 where
176 F: Fn(Vec<String>) -> Fut + Send + Sync + 'static,
177 Fut: Future<Output = ()> + Send + 'static,
178 {
179 self.callbacks.on_tool_cancelled = Some(Arc::new(move |ids| Box::pin(f(ids))));
180 self
181 }
182
183 pub fn on_turn_complete<F, Fut>(mut self, f: F) -> Self
185 where
186 F: Fn() -> Fut + Send + Sync + 'static,
187 Fut: Future<Output = ()> + Send + 'static,
188 {
189 self.callbacks.on_turn_complete = Some(Arc::new(move || Box::pin(f())));
190 self
191 }
192
193 pub fn on_generation_complete<F, Fut>(mut self, f: F) -> Self
200 where
201 F: Fn() -> Fut + Send + Sync + 'static,
202 Fut: Future<Output = ()> + Send + 'static,
203 {
204 self.callbacks.on_generation_complete = Some(Arc::new(move || Box::pin(f())));
205 self
206 }
207
208 pub fn on_go_away<F, Fut>(mut self, f: F) -> Self
210 where
211 F: Fn(Duration) -> Fut + Send + Sync + 'static,
212 Fut: Future<Output = ()> + Send + 'static,
213 {
214 self.callbacks.on_go_away = Some(Arc::new(move |d| Box::pin(f(d))));
215 self
216 }
217
218 pub fn on_connected<F, Fut>(mut self, f: F) -> Self
222 where
223 F: Fn(Arc<dyn gemini_genai_rs::session::SessionWriter>) -> Fut + Send + Sync + 'static,
224 Fut: Future<Output = ()> + Send + 'static,
225 {
226 self.callbacks.on_connected = Some(Arc::new(move |w| Box::pin(f(w))));
227 self
228 }
229
230 pub fn on_disconnected<F, Fut>(mut self, f: F) -> Self
232 where
233 F: Fn(Option<String>) -> Fut + Send + Sync + 'static,
234 Fut: Future<Output = ()> + Send + 'static,
235 {
236 self.callbacks.on_disconnected = Some(Arc::new(move |r| Box::pin(f(r))));
237 self
238 }
239
240 pub fn on_teardown<F, Fut>(mut self, f: F) -> Self
259 where
260 F: Fn() -> Fut + Send + Sync + 'static,
261 Fut: Future<Output = ()> + Send + 'static,
262 {
263 self.callbacks
264 .on_teardown
265 .push(Arc::new(move || Box::pin(f())));
266 self
267 }
268
269 pub fn on_resumed<F, Fut>(mut self, f: F) -> Self
274 where
275 F: Fn() -> Fut + Send + Sync + 'static,
276 Fut: Future<Output = ()> + Send + 'static,
277 {
278 self.callbacks.on_resumed = Some(Arc::new(move || Box::pin(f())));
279 self
280 }
281
282 pub fn on_error<F, Fut>(mut self, f: F) -> Self
284 where
285 F: Fn(String) -> Fut + Send + Sync + 'static,
286 Fut: Future<Output = ()> + Send + 'static,
287 {
288 self.callbacks.on_error = Some(Arc::new(move |e| Box::pin(f(e))));
289 self
290 }
291
292 pub fn on_turn_complete_concurrent<F, Fut>(mut self, f: F) -> Self
298 where
299 F: Fn() -> Fut + Send + Sync + 'static,
300 Fut: Future<Output = ()> + Send + 'static,
301 {
302 self.callbacks.on_turn_complete = Some(Arc::new(move || Box::pin(f())));
303 self.callbacks.on_turn_complete_mode = CallbackMode::Concurrent;
304 self
305 }
306
307 pub fn on_generation_complete_concurrent<F, Fut>(mut self, f: F) -> Self
309 where
310 F: Fn() -> Fut + Send + Sync + 'static,
311 Fut: Future<Output = ()> + Send + 'static,
312 {
313 self.callbacks.on_generation_complete = Some(Arc::new(move || Box::pin(f())));
314 self.callbacks.on_generation_complete_mode = CallbackMode::Concurrent;
315 self
316 }
317
318 pub fn on_connected_concurrent<F, Fut>(mut self, f: F) -> Self
320 where
321 F: Fn(Arc<dyn gemini_genai_rs::session::SessionWriter>) -> Fut + Send + Sync + 'static,
322 Fut: Future<Output = ()> + Send + 'static,
323 {
324 self.callbacks.on_connected = Some(Arc::new(move |w| Box::pin(f(w))));
325 self.callbacks.on_connected_mode = CallbackMode::Concurrent;
326 self
327 }
328
329 pub fn on_disconnected_concurrent<F, Fut>(mut self, f: F) -> Self
331 where
332 F: Fn(Option<String>) -> Fut + Send + Sync + 'static,
333 Fut: Future<Output = ()> + Send + 'static,
334 {
335 self.callbacks.on_disconnected = Some(Arc::new(move |r| Box::pin(f(r))));
336 self.callbacks.on_disconnected_mode = CallbackMode::Concurrent;
337 self
338 }
339
340 pub fn on_resumed_concurrent<F, Fut>(mut self, f: F) -> Self
342 where
343 F: Fn() -> Fut + Send + Sync + 'static,
344 Fut: Future<Output = ()> + Send + 'static,
345 {
346 self.callbacks.on_resumed = Some(Arc::new(move || Box::pin(f())));
347 self.callbacks.on_resumed_mode = CallbackMode::Concurrent;
348 self
349 }
350
351 pub fn on_error_concurrent<F, Fut>(mut self, f: F) -> Self
353 where
354 F: Fn(String) -> Fut + Send + Sync + 'static,
355 Fut: Future<Output = ()> + Send + 'static,
356 {
357 self.callbacks.on_error = Some(Arc::new(move |e| Box::pin(f(e))));
358 self.callbacks.on_error_mode = CallbackMode::Concurrent;
359 self
360 }
361
362 pub fn on_go_away_concurrent<F, Fut>(mut self, f: F) -> Self
364 where
365 F: Fn(Duration) -> Fut + Send + Sync + 'static,
366 Fut: Future<Output = ()> + Send + 'static,
367 {
368 self.callbacks.on_go_away = Some(Arc::new(move |d| Box::pin(f(d))));
369 self.callbacks.on_go_away_mode = CallbackMode::Concurrent;
370 self
371 }
372
373 pub fn on_tool_cancelled_concurrent<F, Fut>(mut self, f: F) -> Self
375 where
376 F: Fn(Vec<String>) -> Fut + Send + Sync + 'static,
377 Fut: Future<Output = ()> + Send + 'static,
378 {
379 self.callbacks.on_tool_cancelled = Some(Arc::new(move |ids| Box::pin(f(ids))));
380 self.callbacks.on_tool_cancelled_mode = CallbackMode::Concurrent;
381 self
382 }
383
384 pub fn on_extracted_concurrent<F, Fut>(mut self, f: F) -> Self
386 where
387 F: Fn(String, serde_json::Value) -> Fut + Send + Sync + 'static,
388 Fut: Future<Output = ()> + Send + 'static,
389 {
390 self.callbacks.on_extracted = Some(Arc::new(move |name, value| Box::pin(f(name, value))));
391 self.callbacks.on_extracted_mode = CallbackMode::Concurrent;
392 self
393 }
394
395 pub fn on_extraction_error_concurrent<F, Fut>(mut self, f: F) -> Self
397 where
398 F: Fn(String, String) -> Fut + Send + Sync + 'static,
399 Fut: Future<Output = ()> + Send + 'static,
400 {
401 self.callbacks.on_extraction_error =
402 Some(Arc::new(move |name, error| Box::pin(f(name, error))));
403 self.callbacks.on_extraction_error_mode = CallbackMode::Concurrent;
404 self
405 }
406}
407
408#[cfg(test)]
409mod tests {
410 use super::*;
411
412 #[test]
415 fn builder_accepts_new_callbacks() {
416 let _live = Live::builder()
417 .on_phase(|_phase| {})
419 .on_tool_cancelled(|_ids| async {})
421 .on_generation_complete(|| async {})
423 .on_resumed(|| async {});
425 }
427
428 #[test]
430 fn builder_accepts_new_callbacks_concurrent() {
431 let _live = Live::builder()
432 .on_tool_cancelled_concurrent(|_ids| async {})
433 .on_generation_complete_concurrent(|| async {})
434 .on_resumed_concurrent(|| async {});
435 }
437}