gemini_adk_fluent_rs/live/mod.rs
1//! `Live` — Fluent builder for callback-driven Gemini Live sessions.
2//!
3//! Wraps L1's `LiveSessionBuilder` with ergonomic callback registration
4//! and integration with composition modules (M, T, P).
5//!
6//! # Callback Modes
7//!
8//! Control-lane callbacks support two execution modes via [`gemini_adk_rs::live::CallbackMode`]:
9//!
10//! - **Default methods** (e.g., `.on_turn_complete()`) → [`gemini_adk_rs::live::CallbackMode::Blocking`]
11//! - **`_concurrent` methods** (e.g., `.on_turn_complete_concurrent()`) → [`gemini_adk_rs::live::CallbackMode::Concurrent`]
12//!
13//! Use concurrent mode for fire-and-forget work (logging, analytics, webhook dispatch).
14//!
15//! # Background Tool Execution
16//!
17//! Mark tools for background execution to eliminate dead air in voice sessions:
18//!
19//! ```rust,ignore
20//! Live::builder()
21//! .tools(dispatcher)
22//! .tool_background("search_kb")
23//! .connect_vertex(project, location, token)
24//! .await?;
25//! ```
26
27mod callbacks;
28mod config;
29mod connect;
30mod contract;
31mod extraction;
32mod introspect;
33
34/// The ambient-tool merge `connect` performs, exposed so
35/// [`check_live`](crate::testing::check_live) checks the same flow the session
36/// will actually run rather than the one the caller wrote.
37pub(crate) use connect::merge_ambient as merge_ambient_for_check;
38mod phases;
39
40use std::collections::HashMap;
41use std::sync::Arc;
42use std::time::Duration;
43
44pub use gemini_adk_rs::live::extractor::TurnExtractor;
45pub use gemini_adk_rs::live::needs::RepairConfig;
46pub use gemini_adk_rs::live::persistence::SessionPersistence;
47pub use gemini_adk_rs::live::steering::{ContextDelivery, SteeringMode};
48pub use gemini_adk_rs::live::{
49 ComputedRegistry, EventCallbacks, InstructionModifier, Phase, TemporalRegistry,
50 ToolExecutionMode, WatcherRegistry,
51};
52use gemini_adk_rs::llm::BaseLlm;
53use gemini_adk_rs::tool::ToolDispatcher;
54use gemini_adk_rs::State;
55use gemini_genai_rs::prelude::*;
56
57// Carve (gap #9): `gemini_adk_fluent_rs::live` is the curated home for the full
58// Live control plane. The kernel `prelude` keeps only `Live` + the headline types;
59// everything else (persistence, steering, repair, transcripts, extraction triggers,
60// soft-turn, runtime contract, …) is re-exported here. (Explicit, rather than a
61// glob, to avoid shadowing the L1/L2 private `callbacks`/`contract` modules.)
62pub use gemini_adk_rs::live::{
63 BackendInputVad, BackendVadSnapshot, BackgroundAgentDispatcher, BackgroundToolTracker,
64 CallbackMode, ComputedContract, ComputedVar, ConsecutiveFailureDetector, ContextBuilder,
65 ControlContract, DefaultResultFormatter, DeferredWriter, EffectMode, EffectPolicy,
66 ExtractionTrigger, ExtractorContract, FieldPromotion, FsPersistence, LiveEffect,
67 LiveEffectExecutor, LiveEvent, LiveEventStream, LiveHandle, LiveReactor, LiveSessionBuilder,
68 LlmExtractor, MemoryPersistence, MergePolicy, NeedsFulfillment, PatternDetector,
69 PendingContext, PhaseContract, PhaseInstruction, PhaseMachine, PhasePreparation,
70 PhaseTransition, PredicateFn, PreparationContract, PromotionContract, RateDetector, Reaction,
71 ReactorEvent, ReactorRule, RepairAction, ResultFormatter, RuntimeContract, SessionSignals,
72 SessionSnapshot, SessionTelemetry, SessionType, SoftTurnDetector, SustainedDetector,
73 ToolCallSummary, ToolContract, TranscriptBuffer, TranscriptTurn, TranscriptWindow, Transition,
74 TransitionContract, TransitionEvaluation, TransitionResult, TransitionTrigger,
75 TurnCountDetector, VoiceRuntimeState, WatchPredicate, Watcher, WatcherContract,
76};
77// Offline record/replay harness (Milestone 7 determinism spine).
78pub use gemini_adk_rs::live::replay::{
79 attach_session, collect_events_until_idle, replay_session, ReplaySession,
80};
81
82/// A deferred agent tool registration (resolved at connect time when State is available).
83pub(crate) struct DeferredAgentTool {
84 pub(crate) name: String,
85 pub(crate) description: String,
86 pub(crate) agent: Arc<dyn gemini_adk_rs::text::TextAgent>,
87}
88
89/// Fluent builder for constructing and connecting Gemini Live sessions.
90///
91/// Accumulates model configuration, callbacks, extractors, phases, watchers,
92/// temporal patterns, and tool execution modes, then connects via one of
93/// the `connect_*` methods.
94///
95/// Control-lane callbacks can be registered with `_concurrent` suffixed
96/// methods for fire-and-forget execution. Tools can be marked for background
97/// execution via [`tool_background()`](Self::tool_background).
98///
99/// # Example
100/// ```ignore
101/// let session = Live::builder()
102/// .model(GeminiModel::Gemini2_0FlashLive)
103/// .voice(Voice::Kore)
104/// .instruction("You are a weather assistant")
105/// .tools(dispatcher)
106/// .on_audio(|data| playback_tx.send(data.clone()).ok())
107/// .on_text(|t| print!("{t}"))
108/// .on_interrupted(|| async { playback.flush().await; })
109/// .connect_vertex("project", "us-central1", token)
110/// .await?;
111/// ```
112///
113/// # Extraction Pipeline
114/// ```ignore
115/// let handle = Live::builder()
116/// .model(GeminiModel::Gemini2_0FlashLive)
117/// .instruction("You are a restaurant order assistant")
118/// .extract_turns::<OrderState>(
119/// flash_llm,
120/// "Extract: items ordered, quantities, modifications, order_phase",
121/// )
122/// .on_extracted(|name, value| async move {
123/// println!("Extracted {name}: {value}");
124/// })
125/// .connect_vertex(project, location, token)
126/// .await?;
127///
128/// // Read latest extraction from shared State at any time:
129/// let order: Option<OrderState> = handle.extracted("OrderState");
130/// ```
131pub struct Live {
132 pub(crate) config: SessionConfig,
133 pub(crate) callbacks: EventCallbacks,
134 pub(crate) dispatcher: Option<ToolDispatcher>,
135 pub(crate) extractors: Vec<Arc<dyn TurnExtractor>>,
136 // L1 registries
137 pub(crate) computed: ComputedRegistry,
138 pub(crate) phases: Vec<Phase>,
139 pub(crate) initial_phase: Option<String>,
140 pub(crate) watchers: WatcherRegistry,
141 pub(crate) temporal: TemporalRegistry,
142 pub(crate) greeting: Option<String>,
143 // Phase defaults: modifiers + prompt_on_enter inherited by all phases.
144 pub(crate) phase_default_modifiers: Vec<InstructionModifier>,
145 pub(crate) phase_default_prompt_on_enter: bool,
146 // Per-tool execution modes (standard vs background).
147 pub(crate) tool_execution_modes: HashMap<String, ToolExecutionMode>,
148 // Deferred agent tools (resolved at connect time).
149 pub(crate) deferred_agent_tools: Vec<DeferredAgentTool>,
150 // Tools requiring async I/O to resolve (MCP/A2A/OpenAPI/Search),
151 // resolved at connect time.
152 pub(crate) deferred_tools: Vec<crate::compose::tools::DeferredTool>,
153 // LLMs to warm up at connect time.
154 pub(crate) warm_up_llms: Vec<Arc<dyn BaseLlm>>,
155 // Control plane configuration.
156 pub(crate) soft_turn_timeout: Option<Duration>,
157 pub(crate) steering_mode: SteeringMode,
158 pub(crate) context_delivery: ContextDelivery,
159 pub(crate) delivery: gemini_adk_rs::live::DeliveryConfig,
160 pub(crate) repair_config: Option<RepairConfig>,
161 pub(crate) persistence: Option<Arc<dyn SessionPersistence>>,
162 pub(crate) session_id: Option<String>,
163 pub(crate) tool_advisory: bool,
164 pub(crate) telemetry_interval: Option<Duration>,
165 // Middleware layers run around tool dispatch in the control lane.
166 pub(crate) middleware_layers: Vec<Arc<dyn gemini_adk_rs::middleware::Middleware>>,
167 // Confirmation provider consulted before running `T::confirm(..)` tools.
168 pub(crate) confirmation_provider:
169 Option<Arc<dyn gemini_adk_rs::confirmation::ConfirmationProvider>>,
170 // Governed flow (DAG) + its enforcement mode.
171 pub(crate) flow: Option<gemini_adk_rs::flow::Flow>,
172 pub(crate) flow_mode: gemini_adk_rs::flow::Enforcement,
173 /// Merged into the flow's own `ambient` list at connect, so an extension
174 /// that registers cross-cutting tools composes with `govern` in either order.
175 pub(crate) ambient_tools: Vec<String>,
176 /// Set by `govern_compiled`/`observe_compiled`, whose documented contract is
177 /// that a `CompiledFlow` already surfaced its diagnostics and connect will
178 /// not re-check it. Connect validates the flow only when this is false.
179 pub(crate) flow_precompiled: bool,
180 /// Caller-supplied session `State`, so tools and flow guards can share one.
181 pub(crate) state: Option<State>,
182 // Per-step on_enter actions: run an agent in a mode when a step activates.
183 pub(crate) flow_actions: Vec<(
184 String,
185 Arc<dyn gemini_adk_rs::text::TextAgent>,
186 gemini_adk_rs::orchestration::Mode,
187 )>,
188 // Wire-log path: a FileWireRecorder is created here at connect time.
189 pub(crate) record_wire_path: Option<std::path::PathBuf>,
190}
191
192impl Live {
193 /// Start building a Live session.
194 ///
195 /// # Examples
196 ///
197 /// Minimal live session setup:
198 ///
199 /// ```rust,ignore
200 /// use gemini_adk_fluent_rs::prelude::*;
201 ///
202 /// let handle = Live::builder()
203 /// .model(GeminiModel::Gemini2_0FlashLive)
204 /// .voice(Voice::Kore)
205 /// .instruction("You are a helpful assistant")
206 /// .greeting("Hello! How can I help?")
207 /// .on_audio(|data| { /* send to speaker */ })
208 /// .on_text(|t| print!("{t}"))
209 /// .connect_google_ai("API_KEY")
210 /// .await?;
211 ///
212 /// handle.send_text("What is the weather?").await?;
213 /// handle.disconnect().await?;
214 /// ```
215 ///
216 /// With phases and state-based transitions:
217 ///
218 /// ```rust,ignore
219 /// let handle = Live::builder()
220 /// .model(GeminiModel::Gemini2_0FlashLive)
221 /// .phase("greeting")
222 /// .instruction("Welcome the user")
223 /// .transition("main", S::is_true("greeted"))
224 /// .done()
225 /// .phase("main")
226 /// .instruction("Help the user")
227 /// .terminal()
228 /// .done()
229 /// .initial_phase("greeting")
230 /// .connect_google_ai("API_KEY")
231 /// .await?;
232 /// ```
233 pub fn builder() -> Self {
234 Self {
235 config: SessionConfig::from_endpoint(ApiEndpoint::google_ai("")),
236 callbacks: EventCallbacks::default(),
237 dispatcher: None,
238 extractors: Vec::new(),
239 computed: ComputedRegistry::new(),
240 phases: Vec::new(),
241 initial_phase: None,
242 watchers: WatcherRegistry::new(),
243 temporal: TemporalRegistry::new(),
244 greeting: None,
245 phase_default_modifiers: Vec::new(),
246 phase_default_prompt_on_enter: false,
247 tool_execution_modes: HashMap::new(),
248 deferred_agent_tools: Vec::new(),
249 deferred_tools: Vec::new(),
250 warm_up_llms: Vec::new(),
251 soft_turn_timeout: None,
252 steering_mode: SteeringMode::default(),
253 context_delivery: ContextDelivery::default(),
254 delivery: gemini_adk_rs::live::DeliveryConfig::default(),
255 repair_config: None,
256 persistence: None,
257 session_id: None,
258 tool_advisory: true,
259 telemetry_interval: None,
260 middleware_layers: Vec::new(),
261 confirmation_provider: None,
262 flow: None,
263 flow_mode: gemini_adk_rs::flow::Enforcement::Enforce,
264 ambient_tools: Vec::new(),
265 flow_precompiled: false,
266 state: None,
267 flow_actions: Vec::new(),
268 record_wire_path: None,
269 }
270 }
271
272 /// Govern the session with a [`Flow`](gemini_adk_rs::flow::Flow) DAG and
273 /// **enforce** it: inadmissible tool calls are blocked and active-step
274 /// postures steer the model at each turn boundary.
275 pub fn govern(mut self, flow: gemini_adk_rs::flow::Flow) -> Self {
276 self.flow = Some(flow);
277 self.flow_mode = gemini_adk_rs::flow::Enforcement::Enforce;
278 self.flow_precompiled = false;
279 self
280 }
281
282 /// Use a `State` you already hold as the session's state.
283 ///
284 /// Without this a tool closure captures whatever `State` the caller made,
285 /// the session runs on a different one, and the two never meet — so a tool
286 /// that writes `identity_verified` and a `Guard::is_true("identity_verified")`
287 /// that reads it are talking about different maps. The guard never fires,
288 /// the flow never advances, and every subsequent tool is refused by a gate
289 /// whose condition was in fact satisfied.
290 ///
291 /// That is the ordinary shape of a governed flow — tools write the facts,
292 /// guards read them — so this is how you make it work:
293 ///
294 /// ```no_run
295 /// # use gemini_adk_fluent_rs::live::Live;
296 /// # use gemini_adk_rs::State;
297 /// let state = State::new();
298 /// Live::builder()
299 /// .with_state(state.clone()) // the session runs on this
300 /// .with_tools(my_tools(state)); // and so do the tools
301 /// # fn my_tools(_: State) -> gemini_adk_fluent_rs::compose::tools::ToolComposite { todo!() }
302 /// ```
303 ///
304 /// `agent_tool` already shares state with the agents it wraps; this is the
305 /// same guarantee for ordinary tools.
306 pub fn with_state(mut self, state: State) -> Self {
307 self.state = Some(state);
308 self
309 }
310
311 /// Register cross-cutting tools as
312 /// [ambient](gemini_adk_rs::flow::Flow::ambient): exempt from every step's
313 /// `allow` whitelist, still bound by anything that names them.
314 ///
315 /// Merged into the governing flow at connect, so this composes with
316 /// [`govern`](Self::govern) in **either order**. Without a flow it is inert.
317 ///
318 /// Extensions that install their own tools should call this rather than
319 /// making the application remember to widen every step — `with_memory` does
320 /// exactly that for `recall_context` and `manage_memory`.
321 pub fn ambient_tools<I, S>(mut self, tools: I) -> Self
322 where
323 I: IntoIterator<Item = S>,
324 S: Into<String>,
325 {
326 self.ambient_tools.extend(tools.into_iter().map(Into::into));
327 self
328 }
329
330 /// The cross-cutting tools registered via [`ambient_tools`](Self::ambient_tools).
331 ///
332 /// Introspection for extensions and tests: the flow's own `ambient` list is
333 /// not included, because the two are only merged at connect.
334 pub fn ambient_tool_names(&self) -> &[String] {
335 &self.ambient_tools
336 }
337
338 /// Attach a [`Flow`](gemini_adk_rs::flow::Flow) in **observe** mode: nothing
339 /// is blocked, but deviations are recorded for audit/analytics.
340 pub fn observe(mut self, flow: gemini_adk_rs::flow::Flow) -> Self {
341 self.flow = Some(flow);
342 self.flow_mode = gemini_adk_rs::flow::Enforcement::Observe;
343 self.flow_precompiled = false;
344 self
345 }
346
347 /// Govern the session with a pre-compiled
348 /// [`CompiledFlow`](gemini_adk_rs::flow::CompiledFlow) and **enforce** it.
349 ///
350 /// A `CompiledFlow` carries proof that
351 /// [`Flow::compile`](gemini_adk_rs::flow::Flow::compile) (or
352 /// [`Flow::compile_with_tools`](gemini_adk_rs::flow::Flow::compile_with_tools))
353 /// already surfaced its diagnostics, so connect does **not** re-validate or
354 /// re-compile it — compile once at load time, govern many sessions.
355 pub fn govern_compiled(self, flow: gemini_adk_rs::flow::CompiledFlow) -> Self {
356 let mut live = self.govern(flow.into_flow());
357 live.flow_precompiled = true;
358 live
359 }
360
361 /// Attach a pre-compiled
362 /// [`CompiledFlow`](gemini_adk_rs::flow::CompiledFlow) in **observe** mode:
363 /// nothing is blocked, but deviations are recorded for audit/analytics.
364 /// Like [`govern_compiled`](Self::govern_compiled), the flow is not
365 /// re-validated or re-compiled at connect.
366 pub fn observe_compiled(self, flow: gemini_adk_rs::flow::CompiledFlow) -> Self {
367 let mut live = self.observe(flow.into_flow());
368 live.flow_precompiled = true;
369 live
370 }
371
372 /// Run an agent the first time the named flow step becomes active.
373 ///
374 /// The agent reads its inputs from `State` and its result lands in
375 /// `{step}:result` ([`AgentMode::Call`] resolves inline at the turn boundary;
376 /// [`AgentMode::Dispatch`]/[`AgentMode::Background`] run detached). A
377 /// downstream step can then complete on it via `Guard::resolved(step)`. This
378 /// is how a governed flow drives in-session orchestration. Requires a flow
379 /// (`govern`/`observe`).
380 ///
381 /// [`AgentMode::Call`]: gemini_adk_rs::orchestration::Mode::Call
382 /// [`AgentMode::Dispatch`]: gemini_adk_rs::orchestration::Mode::Dispatch
383 /// [`AgentMode::Background`]: gemini_adk_rs::orchestration::Mode::Background
384 pub fn on_enter(
385 mut self,
386 step: impl Into<String>,
387 agent: Arc<dyn gemini_adk_rs::text::TextAgent>,
388 mode: gemini_adk_rs::orchestration::Mode,
389 ) -> Self {
390 self.flow_actions.push((step.into(), agent, mode));
391 self
392 }
393
394 /// Gate `T::confirm(..)` tools behind a confirmation provider.
395 ///
396 /// When set, any confirmation-gated tool is checked against `provider`
397 /// before it runs; a denied decision returns an error to the model instead
398 /// of executing the tool. Accepts any [`ConfirmationProvider`] — including a
399 /// plain async closure of `Fn(ConfirmationRequest) -> impl Future<Output = ToolConfirmation>`.
400 ///
401 /// [`ConfirmationProvider`]: gemini_adk_rs::confirmation::ConfirmationProvider
402 /// [`ConfirmationRequest`]: gemini_adk_rs::confirmation::ConfirmationRequest
403 /// [`ToolConfirmation`]: gemini_adk_rs::confirmation::ToolConfirmation
404 pub fn confirmation_provider(
405 mut self,
406 provider: Arc<dyn gemini_adk_rs::confirmation::ConfirmationProvider>,
407 ) -> Self {
408 self.confirmation_provider = Some(provider);
409 self
410 }
411
412 /// Attach a [`MiddlewareComposite`](crate::compose::middleware::MiddlewareComposite)
413 /// — every layer runs around tool
414 /// dispatch in the control lane (`before_tool` can veto a call,
415 /// `after_tool` and `on_tool_error` observe results).
416 ///
417 /// Compose layers with `|`, e.g. `M::log() | M::latency()`.
418 ///
419 /// Note: model-level hooks (`before_model`/`after_model`) are TextAgent
420 /// pipeline concepts and do not apply to a streaming Live session.
421 pub fn middleware(
422 mut self,
423 composite: crate::compose::middleware::MiddlewareComposite,
424 ) -> Self {
425 self.middleware_layers.extend(composite.layers);
426 self
427 }
428
429 /// Set the periodic telemetry emission interval.
430 ///
431 /// When set, the processor emits `LiveEvent::Telemetry` snapshots
432 /// and `LiveEvent::TurnMetrics` at this rate.
433 pub fn telemetry_interval(mut self, interval: Duration) -> Self {
434 self.telemetry_interval = Some(interval);
435 self
436 }
437}
438
439#[cfg(test)]
440mod tests {
441 use super::*;
442 use std::sync::Arc;
443 use std::time::Duration;
444
445 #[test]
446 fn builder_chain_compiles() {
447 let _live = Live::builder()
448 .model(GeminiModel::Gemini2_0FlashLive)
449 .voice(Voice::Kore)
450 .instruction("Test")
451 .temperature(0.7)
452 .google_search()
453 .transcription(true, true)
454 .affective_dialog(true)
455 .session_resume(true)
456 .context_compression(4000, 2000)
457 .on_audio(|_data| {})
458 .on_text(|_t| {})
459 .on_vad_start(|| {})
460 .on_interrupted(|| async {})
461 .on_turn_complete(|| async {})
462 .on_go_away(|_d| async {})
463 .on_connected(|_writer| async {})
464 .on_disconnected(|_r| async {})
465 .on_error(|_e| async {});
466 // Just verify the builder chain compiles
467 }
468
469 #[test]
470 fn govern_compiled_attaches_precompiled_flow_without_recompiling() {
471 use gemini_adk_rs::flow::{Enforcement, Flow, Guard};
472
473 let compiled = Flow::new()
474 .step("greet")
475 .done(Guard::is_true("greeted"))
476 .step("end")
477 .after("greet")
478 .terminal()
479 .build()
480 .expect("valid flow")
481 .compile()
482 .expect("flow compiles");
483
484 // Enforce mode.
485 let live = Live::builder().govern_compiled(compiled.clone());
486 assert!(live.flow.is_some(), "compiled flow attached");
487 assert_eq!(live.flow_mode, Enforcement::Enforce);
488
489 // Observe mode.
490 let live = Live::builder().observe_compiled(compiled);
491 assert!(live.flow.is_some(), "compiled flow attached");
492 assert_eq!(live.flow_mode, Enforcement::Observe);
493 }
494
495 #[test]
496 fn builder_with_extraction_compiles() {
497 use gemini_adk_rs::llm::{BaseLlm, LlmError, LlmRequest, LlmResponse};
498 use schemars::JsonSchema;
499
500 #[derive(serde::Deserialize, serde::Serialize, JsonSchema)]
501 struct OrderState {
502 phase: String,
503 items: Vec<String>,
504 }
505
506 struct FakeLlm;
507
508 #[async_trait::async_trait]
509 impl BaseLlm for FakeLlm {
510 fn model_id(&self) -> &str {
511 "fake"
512 }
513 async fn generate(&self, _req: LlmRequest) -> Result<LlmResponse, LlmError> {
514 unimplemented!()
515 }
516 }
517
518 let _live = Live::builder()
519 .model(GeminiModel::Gemini2_0FlashLive)
520 .instruction("Restaurant order assistant")
521 .extract_turns::<OrderState>(
522 Arc::new(FakeLlm),
523 "Extract order state: items, quantities, phase",
524 )
525 .on_extracted(|name, value| async move {
526 let _ = (name, value);
527 })
528 // Outbound interceptors
529 .before_tool_response(|responses, _state| async move {
530 responses // pass through
531 })
532 .on_turn_boundary(|_state, _writer| async move {
533 // inject context
534 })
535 .instruction_template(|state| {
536 let phase: String = state.get("phase").unwrap_or_default();
537 match phase.as_str() {
538 "ordering" => Some("Take orders accurately.".into()),
539 _ => None,
540 }
541 });
542 // Just verify the builder chain with all features compiles
543 }
544
545 #[test]
546 fn builder_with_computed_state_compiles() {
547 let _live = Live::builder()
548 .model(GeminiModel::Gemini2_0FlashLive)
549 .instruction("Test computed state")
550 .computed("doubled", &["app:count"], |state| {
551 let count: i64 = state.get("app:count")?;
552 Some(serde_json::json!(count * 2))
553 })
554 .computed("level", &["app:score"], |state| {
555 let score: f64 = state.get("app:score")?;
556 if score > 0.5 {
557 Some(serde_json::json!("high"))
558 } else {
559 Some(serde_json::json!("low"))
560 }
561 });
562 }
563
564 #[test]
565 fn builder_with_phases_compiles() {
566 let _live = Live::builder()
567 .model(GeminiModel::Gemini2_0FlashLive)
568 .phase("greeting")
569 .instruction("Welcome the user warmly")
570 .transition("main", |s| s.get::<bool>("greeted").unwrap_or(false))
571 .on_enter(|state, _writer| async move {
572 let _ = state.set("entered_greeting", true);
573 })
574 .done()
575 .phase("main")
576 .dynamic_instruction(|s| {
577 let topic: String = s.get("topic").unwrap_or_default();
578 format!("Discuss {topic}")
579 })
580 .tools(vec!["search".into(), "lookup".into()])
581 .transition("farewell", |s| s.get::<bool>("done").unwrap_or(false))
582 .done()
583 .phase("farewell")
584 .instruction("Say goodbye")
585 .terminal()
586 .done()
587 .initial_phase("greeting");
588 }
589
590 #[test]
591 fn builder_with_phase_guard_compiles() {
592 let _live = Live::builder()
593 .model(GeminiModel::Gemini2_0FlashLive)
594 .phase("start")
595 .instruction("Begin")
596 .transition("secure", |_| true)
597 .done()
598 .phase("secure")
599 .instruction("Secure area")
600 .guard(|s| s.get::<bool>("verified").unwrap_or(false))
601 .on_exit(|state, _writer| async move {
602 let _ = state.set("left_secure", true);
603 })
604 .terminal()
605 .done()
606 .initial_phase("start");
607 }
608
609 #[test]
610 fn builder_with_watchers_compiles() {
611 let _live = Live::builder()
612 .model(GeminiModel::Gemini2_0FlashLive)
613 .watch("app:score")
614 .crossed_above(0.9)
615 .then(|_old, _new, state| async move {
616 let _ = state.set("high_score_alert", true);
617 })
618 .watch("app:status")
619 .changed_to(serde_json::json!("complete"))
620 .blocking()
621 .then(|_old, _new, _state| async move {
622 // blocking action
623 })
624 .watch("app:flag")
625 .became_true()
626 .then(|_old, _new, _state| async move {
627 // flag became true
628 });
629 }
630
631 #[test]
632 fn builder_with_temporal_patterns_compiles() {
633 let _live = Live::builder()
634 .model(GeminiModel::Gemini2_0FlashLive)
635 .when_sustained(
636 "user_confused",
637 |s| s.get::<bool>("confused").unwrap_or(false),
638 Duration::from_secs(30),
639 |_state, _writer| async move {
640 // offer help
641 },
642 )
643 .when_rate(
644 "rapid_errors",
645 |evt| matches!(evt, SessionEvent::TextDelta(_)),
646 5,
647 Duration::from_secs(10),
648 |_state, _writer| async move {
649 // throttle
650 },
651 )
652 .when_turns(
653 "stuck_in_loop",
654 |s| s.get::<bool>("repeating").unwrap_or(false),
655 3,
656 |_state, _writer| async move {
657 // break loop
658 },
659 );
660 }
661
662 #[test]
663 fn builder_full_l1_chain_compiles() {
664 // Full chain combining all L1 features in a single builder
665 let _live = Live::builder()
666 .model(GeminiModel::Gemini2_0FlashLive)
667 .voice(Voice::Kore)
668 .instruction("Full featured agent")
669 // Computed state
670 .computed("sentiment_level", &["app:sentiment_score"], |state| {
671 let score: f64 = state.get("app:sentiment_score")?;
672 if score > 0.7 {
673 Some(serde_json::json!("positive"))
674 } else if score < 0.3 {
675 Some(serde_json::json!("negative"))
676 } else {
677 Some(serde_json::json!("neutral"))
678 }
679 })
680 // Phases
681 .phase("greeting")
682 .instruction("Greet the user")
683 .transition("help", |s| s.get::<bool>("needs_help").unwrap_or(false))
684 .done()
685 .phase("help")
686 .instruction("Help the user")
687 .terminal()
688 .done()
689 .initial_phase("greeting")
690 // Watchers
691 .watch("app:sentiment_score")
692 .crossed_below(0.2)
693 .then(|_old, _new, state| async move {
694 let _ = state.set("alert:low_sentiment", true);
695 })
696 // Temporal
697 .when_turns(
698 "repeated_confusion",
699 |s| s.get::<bool>("confused").unwrap_or(false),
700 3,
701 |_state, _writer| async move {},
702 )
703 // Standard callbacks
704 .on_audio(|_data| {})
705 .on_text(|_t| {})
706 .on_turn_complete(|| async {});
707 }
708
709 #[test]
710 fn builder_with_callback_modes_compiles() {
711 let _live = Live::builder()
712 .model(GeminiModel::Gemini2_0FlashLive)
713 .on_turn_complete_concurrent(|| async {})
714 .on_error_concurrent(|_e| async {})
715 .on_extracted_concurrent(|_name, _val| async {})
716 .on_extraction_error_concurrent(|_name, _err| async {})
717 .on_connected_concurrent(|_w| async {})
718 .on_disconnected_concurrent(|_r| async {})
719 .on_go_away_concurrent(|_d| async {});
720 }
721
722 #[test]
723 fn builder_with_background_tools_compiles() {
724 use gemini_adk_rs::live::DefaultResultFormatter;
725
726 let _live = Live::builder()
727 .model(GeminiModel::Gemini2_0FlashLive)
728 .tool_background("search_kb")
729 .tool_background_with_formatter("analyze_document", Arc::new(DefaultResultFormatter));
730 }
731
732 #[test]
733 fn builder_mixed_callback_modes_and_bg_tools() {
734 use gemini_adk_rs::live::DefaultResultFormatter;
735
736 let _live = Live::builder()
737 .model(GeminiModel::Gemini2_0FlashLive)
738 .voice(Voice::Kore)
739 .instruction("Full featured agent")
740 .tool_background("slow_tool")
741 .tool_background_with_formatter("kb_search", Arc::new(DefaultResultFormatter))
742 .on_turn_complete_concurrent(|| async {})
743 .on_extracted_concurrent(|_name, _val| async {})
744 .on_audio(|_data| {})
745 .on_text(|_t| {})
746 .on_interrupted(|| async {});
747 }
748}