1// 2.2.3 2import { n as qn, T as Wn, J as g, G as Qn, O as fn, V as Xn, A as Zn, S as A, x as V, L as Hn, h as Yn, E as u, z as vn, a as E, b as G, D as ne, c as ee } from "./index-DVPerhZa.js"; 3function te() { 4 try { 5 return Intl.DateTimeFormat().resolvedOptions().timeZone ?? ""; 6 } catch { 7 return ""; 8 } 9} 10function oe() { 11 if (typeof navigator > "u") 12 return ""; 13 const d = navigator.userAgent ?? "", m = (navigator.platform ?? "").toLowerCase(); 14 let s = "Unknown"; 15 /android/i.test(d) ? s = "Android" : /iphone|ipad|ipod/i.test(d) ? s = "iOS" : m.includes("win") ? s = "Windows" : m.includes("mac") ? s = "Mac OS X" : m.includes("linux") && (s = "Linux"); 16 const S = s === "Android" || s === "iOS" || /Mobi/i.test(d); 17 return `${s}, ${S ? "Mobile" : "Desktop"}`; 18} 19function ie() { 20 const d = {}, m = te(); 21 m && (d.timezone = m); 22 const s = oe(); 23 return s && (d.device = s), d; 24} 25function ae(d, m, s, S) { 26 const o = ne(d, `${m}/v2/agents/${s}`, S); 27 return { 28 async createStream(J) { 29 return o.post("/sessions", J); 30 } 31 }; 32} 33const re = { 34 [u.ChatAnswer]: G.Answer, 35 [u.ChatPartial]: G.Partial 36}, mn = 2e4, ce = []; 37async function le() { 38 try { 39 return await import("./index-DVPerhZa.js").then((d) => d.j).then((d) => d.G); 40 } catch { 41 throw new Error( 42 "LiveKit client is required for this streaming manager. Please install it using: npm install livekit-client" 43 ); 44 } 45} 46const se = { 47 excellent: V.Strong, 48 good: V.Strong, 49 poor: V.Weak, 50 lost: V.Unknown, 51 unknown: V.Unknown 52}, O = (d = "Stream Error") => new Xn(d); 53function kn(d, m, s) { 54 var S, o; 55 throw m("Failed to connect to LiveKit room:", d), (S = s.onConnectionStateChange) == null || S.call(s, A.Fail, "internal:init-error"), (o = s.onError) == null || o.call(s, d, { sessionId: "" }), d; 56} 57async function de(d, m, s) { 58 var S; 59 const o = qn(s.debug || !1, "LiveKitStreamingManager"), { Room: J, RoomEvent: h, ConnectionState: $, Track: tn } = await le(), { callbacks: t, auth: Cn, baseURL: bn, analytics: on } = s; 60 let r = null, L = !1; 61 const an = Wn.Fluent; 62 let R = null; 63 const P = { isPublishing: !1, publication: null }, q = { isPublishing: !1, publication: null }; 64 let k = null, y = null, z = null, x = !1; 65 r = new J({ 66 adaptiveStream: !1, 67 // Must be false to use mediaStreamTrack directly 68 dynacast: !0 69 }); 70 for (const [n, e] of s.rpcMethods ?? [])
71 r.registerRpcMethod(n, e); 72 let M = null, I = g.Idle, W = !0; 73 const C = /* @__PURE__ */ new Map(); 74 let F = null, Q = null; 75 const Sn = ae(Cn, bn || Qn, d, t.onError); 76 let b, D, K, rn = !0; 77 try { 78 const n = await Sn.createStream({ 79 transport: m.transport, 80 chat_persist: m.chat_persist ?? !0, 81 verbose: s.verbose ?? !1 82 }), { id: e, session_token: i, session_url: a, interrupt_enabled: c } = n; 83 (S = t.onStreamCreated) == null || S.call(t, { session_id: e, stream_id: e, agent_id: d }), b = e, D = i, K = a, rn = c ?? !0, await r.prepareConnection(K, D); 84 } catch (n) { 85 kn(n, o, t); 86 } 87 if (!K || !D || !b) 88 return Promise.reject(new Error("Failed to initialize LiveKit stream")); 89 r.on(h.ConnectionStateChanged, Tn).on(h.ConnectionQualityChanged, An).on(h.ParticipantConnected, Ln).on(h.ParticipantDisconnected, Mn).on(h.TrackSubscribed, Rn).on(h.TrackUnsubscribed, Pn).on(h.DataReceived, Fn).on(h.MediaDevicesError, jn).on(h.TranscriptionReceived, yn).on(h.EncryptionError, Vn).on(h.TrackSubscriptionFailed, On); 90 function yn(n, e) { 91 e != null && e.isLocal && fn.update(); 92 } 93 async function wn() { 94 const n = ie(); 95 if (!(!r || Object.keys(n).length === 0)) 96 try { 97 await r.localParticipant.setAttributes(n); 98 } catch (e) { 99 o("Failed to set user context attributes", e); 100 } 101 } 102 try { 103 await r.connect(K, D), o("LiveKit room joined successfully"), wn(), M = setTimeout(() => { 104 var n; 105 o( 106 `Track subscription timeout - no track subscribed within ${mn / 1e3} seconds after connect` 107 ), M = null; 108 const e = O("Track subscription timeout"); 109 on.track("connectivity-error", { 110 error: Zn(e), 111 sessionId: b 112 }), (n = t.onError) == null || n.call(t, e, { sessionId: b }), nn("internal:track-subscription-timeout"); 113 }, mn); 114 } catch (n) { 115 kn(n, o, t); 116 } 117 on.enrich({ 118 "stream-type": an 119 }); 120 function Tn(n) { 121 var e, i, a, c; 122 switch (o("Connection state changed:", n), n) { 123 case $.Connecting: 124 o("CALLBACK: onConnectionStateChange(Connecting)"), (e = t.onConnectionStateChange) == null || e.call(t, A.Connecting, "livekit:connecting"); 125 break; 126 case $.Connected: 127 o("LiveKit room connected successfully"), L = !0; 128 break; 129 case $.Disconnected: 130 o("LiveKit room disconnected"), L = !1, x = !1, P.publication = null, q.publication = null, (i = t.onConnectionStateChange) == null || i.call( 131 t, 132 A.Disconnected, 133 Q ?? "livekit:disconnected" 134 ); 135 break; 136 case $.Reconnecting: 137 o("LiveKit room reconnecting..."), (a = t.onConnectionStateChange) == null || a.call(t, A.Connecting, "livekit:reconnecting"); 138 break; 139 case $.SignalReconnecting: 140 o("LiveKit room signal reconnecting..."), (c = t.onConnectionStateChange) == null || c.call(t, A.Connecting, "livekit:signal-reconnecting"); 141 break; 142 } 143 } 144 function An(n, e) { 145 var i; 146 o("Connection quality:", n), e != null && e.isLocal && ((i = t.onConnectivityStateChange) == null || i.call(t, se[n])); 147 } 148 function Ln(n) { 149 o("Participant connected:", n.identity); 150 } 151 function Mn(n) { 152 o("Participant disconnected:", n.identity), nn("livekit:participant-disconnected"); 153 } 154 function En() { 155 var n; 156 z !== E.Start && (o("CALLBACK: onVideoStateChange(Start)"), z = E.Start, (n = t.onVideoStateChange) == null || n.call(t, E.Start)); 157 } 158 function cn(n) { 159 var e; 160 z !== E.Stop && (o("CALLBACK: onVideoStateChange(Stop)"), z = E.Stop, (e = t.onVideoStateChange) == null || e.call(t, E.Stop, n)); 161 } 162 function Rn(n, e, i) { 163 var a, c, p; 164 o(`Track subscribed: ${n.kind} from ${i.identity}`); 165 const l = n.mediaStreamTrack; 166 if (!l) { 167 o(`No mediaStreamTrack available for ${n.kind}`); 168 return; 169 } 170 R ? (R.addTrack(l), o(`Added ${n.kind} track to shared MediaStream`)) : (R = new MediaStream([l]), o(`Created shared MediaStream with ${n.kind} track`)), n.kind === "audio" && (y = Hn( 171 () => n.getRTCStatsReport(), 172 ({ sttLatency: f, serviceLatency: v }) => { 173 var j, w, _; 174 const T = fn.get(!0); 175 let gn = 0; 176 if (f) { 177 const hn = ((w = (j = k == null ? void 0 : k.getReport()) == null ? void 0 : j.webRTCStats) == null ? void 0 : w.avgRtt) ?? 0; 178 gn = hn > 0 ? Math.round(hn * 1e3) : 0; 179 } 180 const en = T > 0 ? T + (f ?? 0) + gn : void 0, Jn = en !== void 0 && v !== void 0 ? en - v : void 0; 181 (_ = t.onFirstAudioDetected) == null || _.call(t, { latency: en, networkLatency: Jn }); 182 } 183 )), n.kind === "video" && ((a = t.onStreamReady) == null || a.call(t), o("CALLBACK: onSrcObjectReady"), (c = t.onSrcObjectReady) == null || c.call(t, R), x || (x = !0, o("CALLBACK: onConnectionStateChange(Connected)"), (p = t.onConnectionStateChange) == null || p.call(t, A.Connected, "livekit:track-subscribed")), k = Yn( 184 () => n.getRTCStatsReport(), 185 () => L, 186 ee, 187 (f, v) => { 188 o(`Video state change: ${f}
188`), f === E.Start ? (M && (clearTimeout(M), M = null, o("Track subscription timeout cleared")), En()) : f === E.Stop && cn(v); 189 } 190 ), k.start()); 191 } 192 function Pn(n, e, i) { 193 o(`Track unsubscribed: ${n.kind} from ${i.identity}`), n.kind === "audio" && (y == null || y.destroy(), y = null), n.kind === "video" && (cn((k == null ? void 0 : k.getReport()) ?? void 0), k == null || k.stop(), k = null); 194 } 195 function ln(n, e) { 196 var i; 197 const a = re[n]; 198 a && ((i = t.onMessage) == null || i.call(t, a, { event: a, ...e })); 199 } 200 function X(n, e) { 201 var i, a, c, p; 202 if (n === u.ToolCallStarted) { 203 const l = e; 204 C.set(l.call_id, { 205 call: { 206 callId: l.call_id, 207 name: l.name, 208 executionMode: l.execution_mode === "async" ? "async" : "blocking" 209 }, 210 interruptible: l.interruptible === !0, 211 turnId: l.turn_id ?? F 212 }), N(), U(), I = g.ToolActive, (i = t.onAgentActivityStateChange) == null || i.call(t, g.ToolActive), (a = t.onToolEvent) == null || a.call(t, u.ToolCallStarted, l); 213 return; 214 } 215 if (n === u.ToolCallDone) { 216 const l = e; 217 sn(l.call_id), (c = t.onToolEvent) == null || c.call(t, u.ToolCallDone, l); 218 return; 219 } 220 if (n === u.ToolCallError) { 221 const l = e; 222 sn(l.call_id), (p = t.onToolEvent) == null || p.call(t, u.ToolCallError, l); 223 } 224 } 225 function sn(n) { 226 C.delete(n) && (N(), U()); 227 } 228 function In(n) { 229 let e = !1; 230 for (const [i, a] of C) 231 a.call.executionMode === "blocking" && (n !== null && a.turnId !== null && a.turnId !== n || (C.delete(i), e = !0)); 232 e && (N(), U()); 233 } 234 function N() { 235 var n; 236 const e = ![...C.values()].some(({ call: i }) => i.executionMode === "blocking"); 237 e !== W && (W = e, (n = t.onInterruptibleChange) == null || n.call(t, e)); 238 } 239 function U() { 240 var n; 241 (n = t.onRunningToolCallsChange) == null || n.call( 242 t, 243 C.size === 0 ? ce : [...C.values()].map(({ call: e }) => e) 244 ); 245 } 246 function $n(n, e) { 247 var i, a, c; 248 if (n === u.StreamVideoCreated) { 249 I = g.Talking, (i = t.onAgentActivityStateChange) == null || i.call(t, g.Talking), y == null || y.arm({ 250 sttLatency: (a = e == null ? void 0 : e.stt) == null ? void 0 : a.latency, 251 serviceLatency: e == null ? void 0 : e.serviceLatency 252 }); 253 return; 254 } 255 if (C.size > 0) { 256 I !== g.ToolActive && (I = g.ToolActive, (c = t.onAgentActivityStateChange) == null || c.call(t, g.ToolActive)); 257 return; 258 } 259 } 260 function B(n, e) { 261 var i, a, c, p; 262 const l = ((a = (i = k == null ? void 0 : k.getReport()) == null ? void 0 : i.webRTCStats) == null ? void 0 : a.avgRtt) ?? 0, f = l > 0 ? Math.round(l / 2 * 1e3) : 0, v = { ...e, downstreamNetworkLatency: f }; 263 s.debug && (c = e == null ? void 0 : e.metadata) != null && c.sentiment && (v.sentiment = { 264 id: e.metadata.sentiment.id, 265 name: e.metadata.sentiment.sentiment 266 }), (p = t.onMessage) == null || p.call(t, n, v), $n(n, e); 267 } 268 function Dn(n, e) { 269 var i; 270 (i = t.onMessage) == null || i.call(t, G.Transcribe, { event: G.Transcribe, ...e }), queueMicrotask(() => { 271 var a; 272 (a = t.onAgentActivityStateChange) == null || a.call(t, g.Loading); 273 }); 274 } 275 function Kn(n, e) { 276 var i; 277 F = (e == null ? void 0 : e.turn_id) ?? null, I = g.Loading, (i = t.onAgentActivityStateChange) == null || i.call(t, g.Loading); 278 } 279 function _n(n, e) { 280 var i; 281 const a = (e == null ? void 0 : e.turn_id) ?? null; 282 F !== null && a !== null && a < F || (In(a), I = g.Idle, (i = t.onAgentActivityStateChange) == null || i.call(t, g.Idle)); 283 } 284 function un(n, e) { 285 Q = (e == null ? void 0 : e.reason) ?? null; 286 } 287 const xn = { 288 [u.ChatAnswer]: ln, 289 [u.ChatPartial]: ln, 290 [u.ToolCallStarted]: X, 291 [u.ToolCallDone]: X, 292 [u.ToolCallError]: X, 293 [u.StreamVideoCreated]: B, 294 [u.StreamVideoDone]: B, 295 [u.StreamVideoError]: B, 296 [u.StreamVideoRejected]: B, 297 [u.ChatAudioTranscribed]: Dn, 298 [u.TurnStarted]: Kn, 299 [u.TurnEnded]: _n, 300 [u.StreamDone]: un, 301 [u.StreamFailed]: un 302 }; 303 function Fn(n, e, i, a) { 304 const c = new TextDecoder().decode(n); 305 let p; 306 try { 307 p = JSON.parse(c); 308 } catch (v) { 309 o("Failed to parse data channel message:", v); 310 return; 311 } 312 const l = a || p.subject; 313 if (o("Data received:", { subject: l, data: p }), !l) return; 314 const f = xn[l]; 315 if (f)
316 try { 317 f(l, p); 318 } catch (v) { 319 console.warn("[LiveKitStreamingManager] Data channel handler failed", { subject: l, error: v }); 320 } 321 } 322 function jn(n) { 323 var e; 324 o("Media devices error:", n), (e = t.onError) == null || e.call(t, O(), { sessionId: b }); 325 } 326 function Vn(n) { 327 var e; 328 o("Encryption error:", n), (e = t.onError) == null || e.call(t, O(), { sessionId: b }); 329 } 330 function On(n, e, i) { 331 o("Track subscription failed:", { trackSid: n, participant: e, reason: i }); 332 } 333 function zn(n, e, i) { 334 for (const [a, c] of i) 335 if (c.source === e && c.track) { 336 const p = c.track.mediaStreamTrack; 337 if (p === n || (p == null ? void 0 : p.id) === n.id) 338 return c; 339 } 340 return null; 341 } 342 async function dn(n, e, i, a, c, p) { 343 var l, f, v; 344 if (!L || !r) 345 throw o(`Room is not connected, cannot publish ${a} stream`), new Error("Room is not connected"); 346 if (n.isPublishing) { 347 o(`${a} publish already in progress, skipping`); 348 return; 349 } 350 const j = i(e); 351 if (j.length === 0) 352 throw new Error(`No ${a} track found in the provided MediaStream`); 353 const w = j[0], _ = zn(w, a, c()); 354 if (_) { 355 o(`${a} track is already published, skipping`, { 356 trackId: w.id, 357 publishedTrackId: (f = (l = _.track) == null ? void 0 : l.mediaStreamTrack) == null ? void 0 : f.id 358 }), n.publication = _; 359 return; 360 } 361 if ((v = n.publication) != null && v.track) { 362 const T = n.publication.track.mediaStreamTrack; 363 T !== w && (T == null ? void 0 : T.id) !== w.id && (o(`Unpublishing existing ${a} track before publishing new one`), await p()); 364 } 365 o(`Publishing ${a} track from provided MediaStream`, { trackId: w.id }), n.isPublishing = !0; 366 try { 367 n.publication = await r.localParticipant.publishTrack(w, { source: a }), o(`${a} track published successfully`, { trackSid: n.publication.trackSid }); 368 } catch (T) { 369 throw o(`Failed to publish ${a} track:`, T), T; 370 } finally { 371 n.isPublishing = !1; 372 } 373 } 374 async function pn(n, e) { 375 if (!(!n.publication || !n.publication.track)) 376 try { 377 r && (await r.localParticipant.unpublishTrack(n.publication.track, !1), o(`${e} track unpublished`)); 378 } catch (i) { 379 o(`Error unpublishing ${e} track:`, i); 380 } finally { 381 n.publication = null; 382 } 383 } 384 async function Nn(n) { 385 return dn( 386 P, 387 n, 388 (e) => e.getAudioTracks(), 389 tn.Source.Microphone, 390 () => r.localParticipant.audioTrackPublications, 391 Z 392 ); 393 } 394 async function Z() { 395 return pn(P, "Microphone"); 396 } 397 async function Un(n) { 398 if (!L || !r) 399 throw o("Cannot replace microphone track: room is not connected"), new Error("Room is not connected"); 400 if (n.kind !== "audio") 401 throw o("Cannot replace microphone track: not an audio track", { kind: n.kind }), new Error("Microphone track must be an audio track"); 402 if (P.isPublishing) 403 throw o("Cannot replace microphone track: publish in progress"), new Error("Microphone publish in progress"); 404 const e = P.publication; 405 if (!e || !e.track) 406 throw o("Cannot replace microphone track: no publication to replace"), new Error("No microphone publication to replace"); 407 try { 408 P.isPublishing = !0, await e.track.replaceTrack(n), o("Microphone track replaced", { trackId: n.id, trackSid: e.trackSid }); 409 } finally { 410 P.isPublishing = !1; 411 } 412 } 413 async function Bn(n) { 414 return dn( 415 q, 416 n, 417 (e) => e.getVideoTracks(), 418 tn.Source.Camera, 419 () => r.localParticipant.videoTrackPublications, 420 H 421 ); 422 } 423 async function H() { 424 return pn(q, "Camera"); 425 } 426 function Gn() { 427 R && (R.getTracks().forEach((n) => n.stop()), R = null); 428 } 429 async function Y(n, e) { 430 var i, a; 431 if (!L || !r) { 432 o("Room is not connected for sending messages"), (i = t.onError) == null || i.call(t, O(), { 433 sessionId: b 434 }); 435 return; 436 } 437 try { 438 await r.localParticipant.sendText(e, { topic: n }), o("Message sent successfully:", e); 439 } catch (c) { 440 o("Failed to send message:", c), (a = t.onError) == null || a.call(t, O(), { sessionId: b }); 441 } 442 } 443 async function nn(n) { 444 var e, i;
445 M && (clearTimeout(M), M = null), y == null || y.destroy(), y = null, r && ((e = t.onConnectionStateChange) == null || e.call(t, A.Disconnecting, n), await Promise.all([Z(), H()]), await r.disconnect()), Gn(), L = !1, x = !1, C.size > 0 && (C.clear(), N(), U()), F = null, (i = t.onAgentActivityStateChange) == null || i.call(t, g.Idle), I = g.Idle; 446 } 447 return { 448 speak(n) { 449 const e = typeof n == "string" ? n : JSON.stringify(n); 450 return Y(vn.Speak, e); 451 }, 452 disconnect: () => nn("user:disconnect"), 453 async reconnect() { 454 var n, e; 455 if ((r == null ? void 0 : r.state) === $.Connected) { 456 o("Room is already connected"); 457 return; 458 } 459 if (!r || !K || !D) 460 throw o("Cannot reconnect: missing room, URL or token"), new Error("Cannot reconnect: session not available"); 461 o("Reconnecting to LiveKit room, state:", r.state), x = !1, Q = null, (n = t.onConnectionStateChange) == null || n.call(t, A.Connecting, "user:reconnect"); 462 try { 463 if (await r.connect(K, D), o("Room reconnected"), L = !0, r.remoteParticipants.size === 0) { 464 if (o("Waiting for agent to join..."), !await new Promise((i) => { 465 const a = setTimeout(() => { 466 r == null || r.off(h.ParticipantConnected, c), i(!1); 467 }, 5e3), c = () => { 468 clearTimeout(a), r == null || r.off(h.ParticipantConnected, c), i(!0); 469 }; 470 r == null || r.on(h.ParticipantConnected, c); 471 })) 472 throw o("Agent did not join within timeout"), await r.disconnect(), new Error("Agent did not rejoin the room"); 473 o("Agent joined, reconnection successful"); 474 } 475 } catch (i) { 476 throw o("Failed to reconnect:", i), (e = t.onConnectionStateChange) == null || e.call(t, A.Fail, "user:reconnect-failed"), i; 477 } 478 }, 479 sendDataChannelMessage: Y, 480 publishMicrophoneStream: Nn, 481 unpublishMicrophoneStream: Z, 482 replaceMicrophoneTrack: Un, 483 publishCameraStream: Bn, 484 unpublishCameraStream: H, 485 interrupt(n) { 486 n !== "text" && Y(vn.Interrupt, ""); 487 }, 488 registerRpcMethod(n, e) {
489 r == null || r.registerRpcMethod(n, e); 490 }, 491 unregisterRpcMethod(n) { 492 r == null || r.unregisterRpcMethod(n); 493 }, 494 sessionId: b, 495 streamId: b, 496 streamType: an, 497 interruptAvailable: rn, 498 isInterruptible: W 499 }; 500} 501export { 502 de as createLiveKitStreamingManager, 503 kn as handleInitError 504};
Line numbers count LF bytes from the start of the resource, as the search results do. Vendor segments are library code the classifier recognised; they are stored but not indexed. Bytes are shown as Latin1 characters, one per byte.