ferrin_core/realtime/
mod.rs1mod session;
14mod tools;
15
16use std::future::IntoFuture;
17use std::pin::Pin;
18use std::sync::Arc;
19use std::task::Context;
20use std::task::Poll;
21use std::time::Duration;
22
23use ferrin_spec::JsonValue;
24use ferrin_spec::RealtimeModelRef;
25use ferrin_spec::realtime_model::ClientSecretOptions;
26use ferrin_tool::ToolSet;
27use futures_core::Stream;
28use tokio::sync::mpsc;
29use tokio::task::JoinSet;
30use tokio_util::sync::CancellationToken;
31
32pub use ferrin_spec::realtime_model::ClientSecret;
33pub use ferrin_spec::realtime_model::ConversationItem;
34pub use ferrin_spec::realtime_model::ConversationRole;
35pub use ferrin_spec::realtime_model::Modality;
36pub use ferrin_spec::realtime_model::RealtimeClientEvent;
37pub use ferrin_spec::realtime_model::RealtimeServerEvent;
38pub use ferrin_spec::realtime_model::RealtimeSessionConfig;
39pub use ferrin_spec::realtime_model::RealtimeToolDefinition;
40pub use ferrin_spec::realtime_model::ResponseCreateOptions;
41pub use ferrin_spec::realtime_model::TranscriptionConfig;
42pub use ferrin_spec::realtime_model::TurnDetection;
43pub use ferrin_spec::realtime_model::TurnDetectionKind;
44pub use session::RealtimeHandle;
45pub use tools::realtime_tool_definitions;
46
47use crate::error::Error;
48use crate::registry::ProviderRegistry;
49use crate::registry::default::resolve_model;
50
51const DEFAULT_EVENT_BUFFER: usize = 256;
53
54const CLOSE_TIMEOUT: Duration = Duration::from_secs(5);
56
57pub type RealtimeEvent = Result<RealtimeServerEvent, Error>;
60
61#[must_use]
66pub fn realtime_session(model: impl Into<RealtimeModelRef>) -> RealtimeSessionBuilder {
67 RealtimeSessionBuilder {
68 model: model.into(),
69 client_secret: None,
70 expires_after_seconds: None,
71 config: RealtimeSessionConfig::default(),
72 tools: ToolSet::new(),
73 tools_context: None,
74 cancellation: CancellationToken::new(),
75 event_buffer: DEFAULT_EVENT_BUFFER,
76 }
77}
78
79#[derive(Debug)]
81pub struct RealtimeSessionBuilder {
82 model: RealtimeModelRef,
83 client_secret: Option<ClientSecret>,
84 expires_after_seconds: Option<u64>,
85 config: RealtimeSessionConfig,
86 tools: ToolSet,
87 tools_context: Option<JsonValue>,
88 cancellation: CancellationToken,
89 event_buffer: usize,
90}
91
92impl RealtimeSessionBuilder {
93 #[must_use]
96 pub fn client_secret(mut self, secret: ClientSecret) -> Self {
97 self.client_secret = Some(secret);
98 self
99 }
100
101 #[must_use]
103 pub fn expires_after_seconds(mut self, seconds: u64) -> Self {
104 self.expires_after_seconds = Some(seconds);
105 self
106 }
107
108 #[must_use]
111 pub fn config(mut self, config: RealtimeSessionConfig) -> Self {
112 self.config = config;
113 self
114 }
115
116 #[must_use]
118 pub fn instructions(mut self, instructions: impl Into<String>) -> Self {
119 self.config.instructions = Some(instructions.into());
120 self
121 }
122
123 #[must_use]
125 pub fn voice(mut self, voice: impl Into<String>) -> Self {
126 self.config.voice = Some(voice.into());
127 self
128 }
129
130 #[must_use]
135 pub fn tools(mut self, tools: ToolSet) -> Self {
136 self.tools = tools;
137 self
138 }
139
140 #[must_use]
142 pub fn tools_context(mut self, context: JsonValue) -> Self {
143 self.tools_context = Some(context);
144 self
145 }
146
147 #[must_use]
149 pub fn cancellation(mut self, cancellation: CancellationToken) -> Self {
150 self.cancellation = cancellation;
151 self
152 }
153
154 #[must_use]
157 pub fn event_buffer(mut self, capacity: usize) -> Self {
158 self.event_buffer = capacity.max(1);
159 self
160 }
161
162 pub async fn connect(self) -> Result<RealtimeSession, Error> {
170 let model = resolve_model(&self.model, ProviderRegistry::realtime_model)?;
171 let mut config = self.config;
172 let definitions =
173 realtime_tool_definitions(&self.tools, self.tools_context.as_ref()).await?;
174 config.tools.extend(definitions);
175
176 let secret = match self.client_secret {
177 Some(secret) => secret,
178 None => model
179 .do_create_client_secret(ClientSecretOptions {
180 expires_after_seconds: self.expires_after_seconds,
181 session_config: Some(config.clone()),
182 })
183 .await
184 .map_err(Error::from)?,
185 };
186
187 let (events_tx, events_rx) = mpsc::channel(self.event_buffer);
188 let mut tasks = JoinSet::new();
189 let handle = session::start(
190 session::StartOptions {
191 model,
192 secret,
193 config,
194 tools: Arc::new(self.tools),
195 tools_context: self.tools_context,
196 cancellation: self.cancellation,
197 events: events_tx,
198 },
199 &mut tasks,
200 )
201 .await?;
202 Ok(RealtimeSession {
203 handle,
204 events: events_rx,
205 tasks,
206 })
207 }
208}
209
210impl IntoFuture for RealtimeSessionBuilder {
211 type Output = Result<RealtimeSession, Error>;
212 type IntoFuture = futures_util::future::BoxFuture<'static, Self::Output>;
213
214 fn into_future(self) -> Self::IntoFuture {
215 Box::pin(self.connect())
216 }
217}
218
219#[derive(Debug)]
226pub struct RealtimeSession {
227 handle: RealtimeHandle,
228 events: mpsc::Receiver<RealtimeEvent>,
229 tasks: JoinSet<()>,
230}
231
232impl RealtimeSession {
233 #[must_use]
235 pub fn handle(&self) -> RealtimeHandle {
236 self.handle.clone()
237 }
238
239 pub async fn send(&self, event: RealtimeClientEvent) -> Result<(), Error> {
245 self.handle.send(event).await
246 }
247
248 pub async fn send_text(&self, text: impl Into<String>) -> Result<(), Error> {
254 self.handle.send_text(text).await
255 }
256
257 pub async fn add_tool_output(&self, call_id: &str, output: &JsonValue) -> Result<(), Error> {
263 self.handle.add_tool_output(call_id, output).await
264 }
265
266 pub async fn next_event(&mut self) -> Option<RealtimeEvent> {
268 self.events.recv().await
269 }
270
271 pub fn events(&mut self) -> impl Stream<Item = RealtimeEvent> + Send + '_ {
273 self
274 }
275
276 #[must_use]
278 pub fn is_closed(&self) -> bool {
279 self.handle.is_closed()
280 }
281
282 pub async fn close(mut self) -> Result<(), Error> {
289 self.handle.close();
290 let started = tokio::time::Instant::now();
291 let deadline = started + CLOSE_TIMEOUT;
292 loop {
293 match tokio::time::timeout_at(deadline, self.tasks.join_next()).await {
294 Ok(None) => return Ok(()),
295 Ok(Some(_)) => {}
296 Err(_) => {
297 self.tasks.abort_all();
298 return Err(Error::Timeout {
299 scope: crate::timeout::TimeoutScope::Total,
300 elapsed: started.elapsed(),
301 });
302 }
303 }
304 }
305 }
306}
307
308impl Stream for RealtimeSession {
309 type Item = RealtimeEvent;
310
311 fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
312 self.events.poll_recv(cx)
313 }
314}