Abonnements en temps réel

Écrit par Stanislas

Dernière mise à jour Il y a 7 mois


S'abonner à des événements en temps réel tels que la diffusion de messages, les appels d'outils et l'activité des utilisateurs à l'aide de connexions WebSocket. Apprenez à configurer des abonnements et à gérer des mises à jour en direct dans votre application de chat.

Présentation

Les abonnements vous permettent d'écouter les événements en temps réel depuis l'API à l'aide de connexions WebSocket. Au lieu d'interroger les mises à jour, votre client reçoit les événements au fur et à mesure qu'ils se produisent. Cela permet d'utiliser des fonctionnalités telles que le streaming de messages caractère par caractère, les indicateurs d'appel d'outils en direct, les notifications de saisie, etc.

L'API Swiftask propose six types d'abonnement principaux : streaming de messages, nouveaux messages complets, messages intermédiaires, appels d'outils d'agent, saisie utilisateur et analyse des sources de données par les bots. Vous pouvez vous abonner à un ou plusieurs événements simultanément.


Prérequis

Avant de configurer les abonnements, assurez-vous que vous disposez des éléments suivants :

  1. Vous êtes authentifié auprès de l'API — Suivez les instructions Guide d'authentification et d'installation pour obtenir votre accessToken et votre workspaceId

  2. Configuré WebSocket dans votre client Apollo — Votre client doit disposer d'un WebSocketLink configuré avec une authentification appropriée

  3. Un identifiant de session — À partir de Gestion des sessions de chat

  4. Compréhension des modèles asynchrones — Les abonnements utilisent des observables et des rappels


Pour commencer

Voici la configuration minimale pour s'abonner aux messages en streaming :

Étape 1 : Configurer la connexion WebSocket

Configurez Apollo Client avec la prise en charge WebSocket pour les abonnements.

import { ApolloClient, InMemoryCache, createHttpLink } from '@apollo/client';
import { setContext } from '@apollo/client/link/context';
import { WebSocketLink } from '@apollo/client/link/ws';
import { split } from '@apollo/client';
import { getMainDefinition } from '@apollo/client/utilities';

const httpLink = createHttpLink({
  uri: 'https://graphql.swiftask.ai/graphql',
});

const authLink = setContext((_, { headers }) => ({
  headers: {
    ...headers,
    authorization: `Bearer ${accessToken}`,
    'x-workspace-id': workspaceId,
    'x-client': 'widget',
  },
}));

const wsLink = new WebSocketLink({
  uri: 'wss://graphql.swiftask.ai/graphql',
  options: {
    reconnect: true,
    connectionParams: {
      authorization: `Bearer ${accessToken}`,
      workspaceId: workspaceId,
    },
  },
});

// Route queries/mutations to HTTP, subscriptions to WebSocket
const splitLink = split(
  ({ query }) => {
    const definition = getMainDefinition(query);
    return definition.kind === 'OperationDefinition' && definition.operation === 'subscription';
  },
  wsLink,
  authLink.concat(httpLink)
);

const client = new ApolloClient({
  link: splitLink,
  cache: new InMemoryCache(),
});

Étape 2 : Abonnez-vous au streaming de messages

Écoutez les blocs de messages en temps réel à mesure qu'ils sont diffusés.

import { gql } from '@apollo/client';

const MESSAGE_STREAM = gql`
  subscription OnMessageStream($sessionId: Float!) {
    onMessageStream(sessionId: $sessionId) {
      messageChunk
      botResponseMessageId
      isStoppable
    }
  }
`;

const subscribeToMessageStream = (client, sessionId, onChunk) => {
  let fullMessage = '';

  return client
    .subscribe({
      query: MESSAGE_STREAM,
      variables: { sessionId },
    })
    .subscribe({
      next: ({ data }) => {
        const chunk = data.onMessageStream.messageChunk;
        fullMessage += chunk;
        onChunk(fullMessage, data.onMessageStream);
      },
      error: (error) => {
        console.error('Stream error:', error);
      },
      complete: () => {
        console.log('Stream complete');
      },
    });
};

// Usage
const subscription = subscribeToMessageStream(client, 67890, (fullMessage, data) => {
  console.log('Streaming:', fullMessage);
  if (data.isStoppable) {
    console.log('User can stop generation');
  }
});

// Cleanup when done
// subscription.unsubscribe();

Étape 3 : Abonnez-vous aux nouveaux messages complets

Écoutez les nouveaux messages complets dans la session (utile pour les chats multi-utilisateurs).

const NEW_MESSAGE = gql`
  subscription NewMessage($sessionId: Float!) {
    newMessage(sessionId: $sessionId) {
      id
      message
      createdAt
      sentBy {
        firstName
      }
      isBotReply
    }
  }
`;

const subscribeToNewMessages = (client, sessionId, onNewMessage) => {
  return client
    .subscribe({
      query: NEW_MESSAGE,
      variables: { sessionId },
    })
    .subscribe({
      next: ({ data }) => {
        onNewMessage(data.newMessage);
      },
      error: (error) => {
        console.error('Subscription error:', error);
      },
    });
};

// Usage
const messageSub = subscribeToNewMessages(client, 67890, (message) => {
  console.log(`${message.sentBy.firstName}: ${message.message}`);
  addMessageToUI(message);
});

Abonnements disponibles

Diffusion de messages

Abonnez-vous au flux en temps réel, caractère par caractère, de la réponse du bot.

Abonnement GraphQL :

subscription OnMessageStream($sessionId: Float!) {
  onMessageStream(sessionId: $sessionId) {
    messageChunk
    userId
    id
    botId
    botResponseMessageId
    isStoppable
  }
}

Champs de réponse :

Champ

Type

Description

messageChunk

chaîne

Le nouveau bloc de texte à ajouter

userId

nombre

ID utilisateur générant le message

id

nombre

ID du message de flux

botId

numéro

ID du bot générant la réponse

botResponseMessageId

numéro

ID du message de réponse complet

isStoppable

booléen

Si la génération peut être arrêtée

Cas d'utilisation :

  • Afficher du texte en streaming en temps réel

  • Afficher un bouton d'arrêt pendant la génération

  • Créer un effet de machine à écrire

  • Mettre à jour l'interface utilisateur caractère par caractère

Exemple avec mises à jour de l'interface utilisateur :

const subscribeToMessageStream = (client, sessionId, onChunk) => {
  let fullMessage = '';
  let messageElement = createMessageElement();

  return client
    .subscribe({
      query: MESSAGE_STREAM,
      variables: { sessionId },
    })
    .subscribe({
      next: ({ data }) => {
        fullMessage += data.onMessageStream.messageChunk;
        messageElement.textContent = fullMessage;

        if (data.isStoppable) {
          showStopButton(() => {
            // Handle stop click
          });
        } else {
          hideStopButton();
        }

        onChunk(fullMessage, data.onMessageStream);
      },
      error: (error) => {
        console.error('Stream error:', error);
        messageElement.classList.add('error');
      },
      complete: () => {
        messageElement.classList.add('complete');
      },
    });
};
New complete messages

Abonnez-vous aux nouveaux messages complets (pas aux fragments diffusés en continu).

Abonnement GraphQL :

subscription NewMessage($sessionId: Float!) {
  newMessage(sessionId: $sessionId) {
    id
    message
    createdAt
    sessionId
    sentBy {
      id
      firstName
      lastName
    }
    isBotReply
  }
}

Cas d'utilisation :

  • Afficher les nouveaux messages dans l'historique des discussions

  • Afficher les notifications pour les nouveaux messages

  • Mettre à jour le nombre de messages

  • Déclencher des animations UI

Exemple :

const NEW_MESSAGE = gql`
  subscription NewMessage($sessionId: Float!) {
    newMessage(sessionId: $sessionId) {
      id
      message
      createdAt
      sentBy {
        firstName
      }
      isBotReply
    }
  }
`;

const subscribeToNewMessages = (client, sessionId, onNewMessage) => {
  return client
    .subscribe({
      query: NEW_MESSAGE,
      variables: { sessionId },
    })
    .subscribe({
      next: ({ data }) => {
        const message = data.newMessage;
        const sender = message.isBotReply ? 'Bot' : message.sentBy.firstName;
        
        onNewMessage(message);
        
        // Play notification sound for user messages
        if (!message.isBotReply) {
          playNotificationSound();
        }
      },
    });
};

Messages intermédiaires

Abonnez-vous aux messages de traitement intermédiaires (les étapes de réflexion du bot).

Abonnement GraphQL :

subscription NewIntermediateMessage($sessionId: Float!) {
  newIntermediateMessage(sessionId: $sessionId) {
    id
    message
    createdAt
    isBotReply
  }
}

Cas d'utilisation :

  • Afficher le processus de réflexion du bot

  • Afficher les étapes d'analyse

  • Déboguer le comportement de l'agent

  • Afficher les indicateurs de progression

Exemple :

const INTERMEDIATE_MESSAGE = gql`
  subscription NewIntermediateMessage($sessionId: Float!) {
    newIntermediateMessage(sessionId: $sessionId) {
      id
      message
      createdAt
    }
  }
`;

const subscribeToIntermediateMessages = (client, sessionId, onIntermediate) => {
  return client
    .subscribe({
      query: INTERMEDIATE_MESSAGE,
      variables: { sessionId },
    })
    .subscribe({
      next: ({ data }) => {
        console.log('Thinking step:', data.newIntermediateMessage.message);
        onIntermediate(data.newIntermediateMessage);
        showThinkingIndicator(data.newIntermediateMessage.message);
      },
    });
};

Appels d'outils de l'agent

Abonnez-vous aux événements lorsque l'agent utilise des outils (recherche, calculs, appels API, etc.).

Abonnement GraphQL :

subscription OnAgentToolCall($sessionId: Float!) {
  onAgentToolCall(sessionId: $sessionId) {
    tool
    toolImage
    botId
    id
    status
  }
}

Champs de réponse :

Champ

Type

Description

tool

chaîne

Nom de l'outil utilisé

toolImage

chaîne

URL de l'icône/image de l'outil

botId

nom

ID du bot utilisant l'outil

id

numéro

ID de l'événement d'appel de l'outil

status

chaîne

« START », « END » ou « ERROR »

Noms d'outils courants :

  • web_search — Recherche sur le Web

  • calculator — Calculs mathématiques

  • knowledge_base — Recherche dans une base de connaissances

  • api_call — Appels API externes

  • database_query — Requêtes de base de données

Cas d'utilisation :

  • Afficher les outils utilisés par l'agent

  • Afficher la progression de l'exécution des outils

  • Afficher les indicateurs d'erreur pour les outils ayant échoué

  • Créer un journal d'activité des outils

Exemple :

const TOOL_CALL = gql`
  subscription OnAgentToolCall($sessionId: Float!) {
    onAgentToolCall(sessionId: $sessionId) {
      tool
      status
      toolImage
    }
  }
`;

const subscribeToToolCalls = (client, sessionId, onToolCall) => {
  const activeTools = {};

  return client
    .subscribe({
      query: TOOL_CALL,
      variables: { sessionId },
    })
    .subscribe({
      next: ({ data }) => {
        const { tool, status, toolImage } = data.onAgentToolCall;

        if (status === 'START') {
          activeTools[tool] = true;
          showToolIndicator(tool, toolImage);
          console.log(`🔧 Using ${tool}...`);
        } else if (status === 'END') {
          delete activeTools[tool];
          hideToolIndicator(tool);
          console.log(`✓ ${tool} completed`);
        } else if (status === 'ERROR') {
          delete activeTools[tool];
          showToolError(tool);
          console.log(`✗ ${tool} failed`);
        }

        onToolCall(data.onAgentToolCall);
      },
    });
};

Saisie de l'utilisateur

S'abonner aux indicateurs de saisie lorsque les utilisateurs tapent.

Abonnement GraphQL :

subscription UserTyping($sessionId: Float!) {
  userTyping(sessionId: $sessionId) {
    id
    firstName
    lastName
  }
}

Cas d'utilisation :

  • Afficher l'indicateur « L'utilisateur est en train de taper... »

  • Améliorer l'expérience de chat en temps réel

  • Empêcher les problèmes de chevauchement des messages

Exemple :

const USER_TYPING = gql`
  subscription UserTyping($sessionId: Float!) {
    userTyping(sessionId: $sessionId) {
      id
      firstName
    }
  }
`;

const subscribeToTyping = (client, sessionId, onTyping) => {
  return client
    .subscribe({
      query: USER_TYPING,
      variables: { sessionId },
    })
    .subscribe({
      next: ({ data }) => {
        const user = data.userTyping;
        showTypingIndicator(`${user.firstName} is typing...`);
        onTyping(user);
      },
    });
};
Bot analyzing datasources

S'abonner aux événements lorsque le bot analyse ses sources de données.

Abonnement GraphQL :

subscription BotAnalyzingDatasources($sessionId: Float!) {
  botAnalyzingDatasources(sessionId: $sessionId) {
    id
    firstName
  }
}

Cas d'utilisation :

  • Afficher l'indicateur « Analyse de la base de connaissances... »

  • Indiquer que le bot recherche des informations

  • Afficher la progression pendant l'analyse des données

Exemple :

const BOT_ANALYZING = gql`
  subscription BotAnalyzingDatasources($sessionId: Float!) {
    botAnalyzingDatasources(sessionId: $sessionId) {
      id
      firstName
    }
  }
`;

const subscribeToAnalyzing = (client, sessionId, onAnalyzing) => {
  return client
    .subscribe({
      query: BOT_ANALYZING,
      variables: { sessionId },
    })
    .subscribe({
      next: ({ data }) => {
        showAnalyzingIndicator('Searching knowledge base...');
        onAnalyzing(data.botAnalyzingDatasources);
      },
    });
};
Complete subscription manager

Voici une classe réutilisable qui gère plusieurs abonnements :

class SubscriptionManager {
  constructor(client, sessionId) {
    this.client = client;
    this.sessionId = sessionId;
    this.subscriptions = [];
  }

  subscribeToAll(callbacks) {
    // Message streaming
    if (callbacks.onMessageChunk) {
      this.subscribeToStream(callbacks.onMessageChunk);
    }

    // New complete messages
    if (callbacks.onNewMessage) {
      this.subscribeToMessages(callbacks.onNewMessage);
    }

    // Intermediate thinking steps
    if (callbacks.onIntermediate) {
      this.subscribeToIntermediateMessages(callbacks.onIntermediate);
    }

    // Tool calls
    if (callbacks.onToolCall) {
      this.subscribeToToolCalls(callbacks.onToolCall);
    }

    // User typing
    if (callbacks.onUserTyping) {
      this.subscribeToTyping(callbacks.onUserTyping);
    }

    // Bot analyzing
    if (callbacks.onBotAnalyzing) {
      this.subscribeToAnalyzing(callbacks.onBotAnalyzing);
    }
  }

  subscribeToStream(onChunk) {
    const MESSAGE_STREAM = gql`
      subscription OnMessageStream($sessionId: Float!) {
        onMessageStream(sessionId: $sessionId) {
          messageChunk
          botResponseMessageId
          isStoppable
        }
      }
    `;

    let fullMessage = '';

    const subscription = this.client
      .subscribe({
        query: MESSAGE_STREAM,
        variables: { sessionId: this.sessionId },
      })
      .subscribe({
        next: ({ data }) => {
          fullMessage += data.onMessageStream.messageChunk;
          onChunk(fullMessage, data.onMessageStream);
        },
      });

    this.subscriptions.push(subscription);
    return subscription;
  }

  subscribeToMessages(onMessage) {
    const NEW_MESSAGE = gql`
      subscription NewMessage($sessionId: Float!) {
        newMessage(sessionId: $sessionId) {
          id
          message
          createdAt
          isBotReply
          sentBy {
            firstName
          }
        }
      }
    `;

    const subscription = this.client
      .subscribe({
        query: NEW_MESSAGE,
        variables: { sessionId: this.sessionId },
      })
      .subscribe({
        next: ({ data }) => {
          onMessage(data.newMessage);
        },
      });

    this.subscriptions.push(subscription);
    return subscription;
  }

  subscribeToIntermediateMessages(onIntermediate) {
    const INTERMEDIATE = gql`
      subscription NewIntermediateMessage($sessionId: Float!) {
        newIntermediateMessage(sessionId: $sessionId) {
          id
          message
          createdAt
        }
      }
    `;

    const subscription = this.client
      .subscribe({
        query: INTERMEDIATE,
        variables: { sessionId: this.sessionId },
      })
      .subscribe({
        next: ({ data }) => {
          onIntermediate(data.newIntermediateMessage);
        },
      });

    this.subscriptions.push(subscription);
    return subscription;
  }

  subscribeToToolCalls(onToolCall) {
    const TOOL_CALL = gql`
      subscription OnAgentToolCall($sessionId: Float!) {
        onAgentToolCall(sessionId: $sessionId) {
          tool
          status
          toolImage
        }
      }
    `;

    const subscription = this.client
      .subscribe({
        query: TOOL_CALL,
        variables: { sessionId: this.sessionId },
      })
      .subscribe({
        next: ({ data }) => {
          onToolCall(data.onAgentToolCall);
        },
      });

    this.subscriptions.push(subscription);
    return subscription;
  }

  subscribeToTyping(onTyping) {
    const USER_TYPING = gql`
      subscription UserTyping($sessionId: Float!) {
        userTyping(sessionId: $sessionId) {
          firstName
        }
      }
    `;

    const subscription = this.client
      .subscribe({
        query: USER_TYPING,
        variables: { sessionId: this.sessionId },
      })
      .subscribe({
        next: ({ data }) => {
          onTyping(data.userTyping);
        },
      });

    this.subscriptions.push(subscription);
    return subscription;
  }

  subscribeToAnalyzing(onAnalyzing) {
    const BOT_ANALYZING = gql`
      subscription BotAnalyzingDatasources($sessionId: Float!) {
        botAnalyzingDatasources(sessionId: $sessionId) {
          id
          firstName
        }
      }
    `;

    const subscription = this.client
      .subscribe({
        query: BOT_ANALYZING,
        variables: { sessionId: this.sessionId },
      })
      .subscribe({
        next: ({ data }) => {
          onAnalyzing(data.botAnalyzingDatasources);
        },
      });

    this.subscriptions.push(subscription);
    return subscription;
  }

  unsubscribeAll() {
    this.subscriptions.forEach((sub) => sub.unsubscribe());
    this.subscriptions = [];
  }

  unsubscribe(subscription) {
    subscription.unsubscribe();
    this.subscriptions = this.subscriptions.filter((s) => s !== subscription);
  }
}

// Usage example
const subManager = new SubscriptionManager(client, 67890);

subManager.subscribeToAll({
  onMessageChunk: (fullMessage, data) => {
    updateStreamingMessage(fullMessage);
  },
  onNewMessage: (message) => {
    addMessageToHistory(message);
  },
  onToolCall: (toolCall) => {
    showToolActivity(toolCall);
  },
  onUserTyping: (user) => {
    showTypingIndicator(user);
  },
  onBotAnalyzing: () => {
    showAnalyzingIndicator();
  },
});

// Cleanup when component unmounts
// subManager.unsubscribeAll();

Cas d'utilisation pratiques

Interface utilisateur de chat complète

Implémentez une interface de chat complète avec tous les types d'abonnement :

class ChatUI {
  constructor(client, sessionId) {
    this.subManager = new SubscriptionManager(client, sessionId);
    this.messages = [];
  }

  initialize() {
    this.subManager.subscribeToAll({
      onMessageChunk: (fullMessage) => {
        this.updateStreamingMessage(fullMessage);
      },
      onNewMessage: (message) => {
        this.messages.push(message);
        this.renderMessages();
        
        if (!message.isBotReply) {
          this.playNotificationSound();
        }
      },
      onToolCall: (toolCall) => {
        if (toolCall.status === 'START') {
          this.showToolIndicator(toolCall.tool);
        } else if (toolCall.status === 'END') {
          this.hideToolIndicator(toolCall.tool);
        } else if (toolCall.status === 'ERROR') {
          this.showToolError(toolCall.tool);
        }
      },
      onUserTyping: (user) => {
        this.showTypingIndicator(`${user.firstName} is typing...`);
      },
      onBotAnalyzing: () => {
        this.showLoadingIndicator('Analyzing knowledge base...');
      },
    });
  }

  updateStreamingMessage(fullMessage) {
    const botMessageElement = document.getElementById('bot-message');
    botMessageElement.textContent = fullMessage;
  }

  renderMessages() {
    const messagesContainer = document.getElementById('messages');
    messagesContainer.innerHTML = this.messages
      .map((msg) => `
        <div class="message ${msg.isBotReply ? 'bot' : 'user'}">
          <p>${msg.message}</p>
          <span class="timestamp">${new Date(msg.createdAt).toLocaleTimeString()}</span>
        </div>
      `)
      .join('');
  }

  showToolIndicator(tool) {
    const indicator = document.createElement('div');
    indicator.id = `tool-${tool}`;
    indicator.className = 'tool-indicator';
    indicator.textContent = `Using ${tool}...`;
    document.getElementById('status').appendChild(indicator);
  }

  hideToolIndicator(tool) {
    const indicator = document.getElementById(`tool-${tool}`);
    if (indicator) indicator.remove();
  }

  showToolError(tool) {
    const indicator = document.getElementById(`tool-${tool}`);
    if (indicator) {
      indicator.classList.add('error');
      indicator.textContent = `${tool} failed`;
    }
  }

  showTypingIndicator(text) {
    document.getElementById('typing').textContent = text;
    setTimeout(() => {
      document.getElementById('typing').textContent = '';
    }, 3000);
  }

  showLoadingIndicator(text) {
    document.getElementById('status').textContent = text;
  }

  playNotificationSound() {
    const audio = new Audio('/notification.mp3');
    audio.play().catch(() => {
      // Audio playback failed, silently continue
    });
  }

  cleanup() {
    this.subManager.unsubscribeAll();
  }
}

// Usage
const chatUI = new ChatUI(client, 67890);
chatUI.initialize();

// Cleanup on page unload
window.addEventListener('beforeunload', () => {
  chatUI.cleanup();
});
Selective subscriptions

Abonnements sélectifs

Abonnez-vous uniquement aux événements dont vous avez besoin :

// Only subscribe to streaming and tool calls
const subManager = new SubscriptionManager(client, sessionId);

subManager.subscribeToStream((fullMessage) => {
  updateUI(fullMessage);
});

subManager.subscribeToToolCalls((toolCall) => {
  if (toolCall.status === 'START') {
    showToolBadge(toolCall.tool);
  }
});

Gestion des erreurs avec les abonnements

Gérez les erreurs WebSocket et les reconnexions :

const subscribeWithErrorHandling = (client, sessionId, onChunk) => {
  const MESSAGE_STREAM = gql`
    subscription OnMessageStream($sessionId: Float!) {
      onMessageStream(sessionId: $sessionId) {
        messageChunk
      }
    }
  `;

  return client
    .subscribe({
      query: MESSAGE_STREAM,
      variables: { sessionId },
    })
    .subscribe({
      next: ({ data }) => {
        onChunk(data.onMessageStream.messageChunk);
      },
      error: (error) => {
        console.error('Subscription error:', error);
        
        if (error.message?.includes('authorization')) {
          // Re-authenticate
          console.log('Re-authenticating...');
          // Handle auth error
        } else if (error.message?.includes('network')) {
          // Network error - retry after delay
          console.log('Network error, retrying in 5s...');
          setTimeout(() => {
            subscribeWithErrorHandling(client, sessionId, onChunk);
          }, 5000);
        }
      },
      complete: () => {
        console.log('Subscription completed');
      },
    });
};

Meilleures pratiques

  1. Désabonnez-vous toujours lorsque vous avez terminé — Évitez les fuites de mémoire en vous désabonnant lorsque les composants sont démontés.

  2. Combinez efficacement les abonnements — Utilisez une classe de gestion pour gérer plusieurs abonnements ensemble.

  3. Gérez les erreurs avec élégance — Implémentez des rappels d'erreur et une logique de reconnexion.

  4. Validez les données avant de mettre à jour l'interface utilisateur — Vérifiez toujours que les données d'abonnement sont valides avant le rendu.

  5. Implémentez des délais d'expiration pour les flux de longue durée — Arrêtez le streaming s'il prend trop de temps.

  6. Affichez les indicateurs d'interface utilisateur appropriés — Utilisez des indicateurs d'outils, des notifications de saisie et des états de chargement.

  7. Mettez en mémoire tampon des blocs pour des mises à jour fluides — Ne mettez pas à jour l'interface utilisateur à chaque caractère ; regroupez les mises à jour pour de meilleures performances.

  8. Utilisez des abonnements sélectifs — N'abonnez-vous qu'aux événements dont vous avez réellement besoin.

  9. Gérez la reconnexion WebSocket — Configurez reconnect: true et implémentez une logique de nouvelle tentative.

  10. Surveillez l'état des abonnements — Enregistrez les événements d'abonnement à des fins de débogage.