PageSourceSearch

https://agent.d-id.com/2.2.3/livekit-manager-25aUaO_G-BExatVlq.js

js d-id.com collected 2026-09-24 08:00:42 UTC 18,806 bytes, 504 lines download raw bytes

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.