All files / api/src/handlers websocket-utils.ts

0% Statements 0/75
100% Branches 1/1
100% Functions 1/1
0% Lines 0/75

Press n or j to go to the next uncovered block, b, p or k for the previous block.

1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93                                                                                                                                                                                         
import type { RealtimeLlmEventParser } from '@api/services/realtime-llm-event-parser';
import type { AppContext } from '@api/types/hono';
import type { InternalProviderAPIConfig } from '@shared/types/ai-providers/config';
import type { SuperAgentsRequestData } from '@shared/types/api/request';
import type { SuperAgentsTarget } from '@shared/types/api/request/headers';
import type { RealtimeSessionOptions } from '@shared/types/realtime';
 
export const addListeners = (
  outgoingWebSocket: WebSocket,
  eventParser: RealtimeLlmEventParser,
  server: WebSocket,
  c: AppContext,
  sessionOptions: RealtimeSessionOptions,
): void => {
  outgoingWebSocket.addEventListener('message', (event) => {
    server?.send(event.data as string);
    try {
      const parsedData = JSON.parse(event.data as string);
      eventParser.handleEvent(c, parsedData, sessionOptions);
    } catch (err: unknown) {
      if (err instanceof Error) {
        console.error('outgoingWebSocket message parse error', err.message);
      } else {
        console.error('outgoingWebSocket message parse error', err);
      }
    }
  });
 
  outgoingWebSocket.addEventListener('close', (event) => {
    server?.close(event.code, event.reason);
  });
 
  outgoingWebSocket.addEventListener('error', (event) => {
    console.error('outgoingWebSocket error', event);
    server?.close();
  });
 
  server.addEventListener('message', (event) => {
    outgoingWebSocket?.send(event.data as string);
  });
 
  server.addEventListener('close', () => {
    outgoingWebSocket?.close();
  });
 
  server.addEventListener('error', (event) => {
    console.error('serverWebSocket error', event);
    outgoingWebSocket?.close();
  });
};
 
export const getOptionsForOutgoingConnection = async (
  c: AppContext,
  apiConfig: InternalProviderAPIConfig,
  saTarget: SuperAgentsTarget,
): Promise<{
  headers: Record<string, string>;
  method: string;
}> => {
  const saRequestData = c.get('sa_request_data');
  const headers = await apiConfig.headers({
    c,
    saTarget,
    saRequestData,
  });
  headers.Upgrade = 'websocket';
  headers.Connection = 'Keep-Alive';
  headers['Keep-Alive'] = 'timeout=600';
  return {
    headers,
    method: 'GET',
  };
};
 
export const getURLForOutgoingConnection = (
  c: AppContext,
  apiConfig: InternalProviderAPIConfig,
  saTarget: SuperAgentsTarget,
  saRequestData: SuperAgentsRequestData,
): string => {
  const baseUrl = apiConfig.getBaseURL({
    c,
    saTarget,
    saRequestData,
  });
  const endpoint = apiConfig.getEndpoint({
    c,
    saTarget,
    saRequestData,
  });
  return `${baseUrl}${endpoint}`;
};