listen method

Future<McpStreamableSubscription> listen({
  1. Object? id,
  2. bool toolsListChanged = false,
  3. bool promptsListChanged = false,
  4. bool resourcesListChanged = false,
  5. Iterable<String> resourceSubscriptions = const <String>[],
  6. Map<String, String> headers = const <String, String>{},
})

Opens a resumable Streamable HTTP listener for server notifications.

Implementation

Future<McpStreamableSubscription> listen({
  Object? id,
  bool toolsListChanged = false,
  bool promptsListChanged = false,
  bool resourcesListChanged = false,
  Iterable<String> resourceSubscriptions = const <String>[],
  Map<String, String> headers = const <String, String>{},
}) async {
  _throwIfClosed();
  if (protocolVersion != latestProtocolVersion) {
    throw const McpStreamableProtocolException(
      'subscriptions/listen requires the MCP 2026 stateless protocol',
    );
  }

  final validatedResources = <String>[];
  final seenResources = <String>{};
  for (final resourceUri in resourceSubscriptions) {
    final uri = _validatedMcpResourceUri(
      resourceUri,
      'resourceSubscriptions',
    );
    if (!seenResources.add(uri)) {
      throw ArgumentError.value(
        resourceSubscriptions,
        'resourceSubscriptions',
        'MCP resource subscriptions must not contain duplicates.',
      );
    }
    validatedResources.add(uri);
  }
  final requestedNotifications = McpSubscriptionFilter(
    toolsListChanged: toolsListChanged,
    promptsListChanged: promptsListChanged,
    resourcesListChanged: resourcesListChanged,
    resourceSubscriptions: validatedResources,
  );
  final requestId = id ?? _nextRequestId++;
  final message = <String, Object?>{
    'jsonrpc': '2.0',
    'id': requestId,
    'method': 'subscriptions/listen',
    'params': <String, Object?>{
      'notifications': requestedNotifications.toJson(),
    },
  };
  _validateJsonRpcRequestId(message, label: 'subscriptions/listen request');
  final preparedMessage = _prepareMessageForProtocol(
    message,
    latestProtocolVersion,
  );
  final requestBody = _encodeBoundedMcpHttpRequest(
    preparedMessage,
    maxRequestBytes,
  );

  HttpClient? pendingSubscriptionHttpClient;
  void abortSubscriptionSetup() {
    final client = pendingSubscriptionHttpClient;
    if (client == null) {
      return;
    }
    pendingSubscriptionHttpClient = null;
    _pendingSubscriptionHttpClients.remove(client);
    client.close(force: true);
  }

  return _runTrackedHttpOperation<McpStreamableSubscription>(
    (operation) async {
      final requestAuthorizationState = _authorizationStateSnapshot;
      final subscriptionStateToken = _subscriptionStateToken;
      final subscriptionHttpClient = _subscriptionHttpClientFactory();
      pendingSubscriptionHttpClient = subscriptionHttpClient;
      _pendingSubscriptionHttpClients.add(subscriptionHttpClient);
      void closeSubscriptionHttpClient() {
        if (identical(
          pendingSubscriptionHttpClient,
          subscriptionHttpClient,
        )) {
          pendingSubscriptionHttpClient = null;
        }
        _pendingSubscriptionHttpClients.remove(subscriptionHttpClient);
        subscriptionHttpClient.close(force: true);
      }

      final HttpClientRequest request;
      try {
        request = await _openTrackedHttpRequest(
          () => subscriptionHttpClient.postUrl(endpoint),
          operation,
          requestAuthorizationState,
          enforceClientState: false,
        );
      } catch (_) {
        closeSubscriptionHttpClient();
        rethrow;
      }
      final HttpClientResponse response;
      try {
        _applyHeaders(
          request,
          accept: _acceptStreamableHttp,
          includeSession: false,
          protocolVersion: latestProtocolVersion,
          authorizationState: requestAuthorizationState,
          extraHeaders: headers,
        );
        _applyStandardRequestHeaders(request, preparedMessage);
        request.headers.contentType = ContentType.json;
        request.persistentConnection = false;
        request.contentLength = requestBody.length;
        request.add(requestBody);
        response = await _sendTrackedHttpRequest(request, operation);
      } catch (_) {
        closeSubscriptionHttpClient();
        rethrow;
      }
      try {
        if (response.statusCode < HttpStatus.ok ||
            response.statusCode >= HttpStatus.multipleChoices) {
          final body = await _readTrackedHttpResponseBody(
            request,
            response,
            operation,
          );
          _throwIfHttpError(response, body);
        }
        if (!_isSse(response)) {
          final body = await _readTrackedHttpResponseBody(
            request,
            response,
            operation,
          );
          if (_isJson(response) && body.isNotEmpty) {
            final jsonResponse = _jsonMapFromBody(
              body,
              'subscriptions/listen JSON response',
            );
            _jsonRpcResultFrom(jsonResponse, method: 'subscriptions/listen');
          }
          throw FormatException(
            'Expected $_acceptSse response, got '
            '${response.headers.contentType?.mimeType ?? 'unknown'}',
          );
        }
        _captureSessionHeaders(
          response,
          captureSessionState: false,
          forbidSessionId: true,
          expectedProtocolVersion: latestProtocolVersion,
        );
      } catch (_) {
        closeSubscriptionHttpClient();
        rethrow;
      }

      if (!identical(subscriptionStateToken, _subscriptionStateToken)) {
        closeSubscriptionHttpClient();
        throw StateError(
          'MCP client closed while subscriptions/listen was pending.',
        );
      }

      late final McpStreamableSubscription subscription;
      subscription = McpStreamableSubscription._(
        id: requestId,
        requestedNotifications: requestedNotifications,
        httpClient: subscriptionHttpClient,
        request: request,
        response: response,
        maxEventBytes: maxResponseBytes,
        onClosed: () => _subscriptions.remove(subscription),
      );
      _pendingSubscriptionHttpClients.remove(subscriptionHttpClient);
      _subscriptions.add(subscription);
      subscription._start();
      try {
        await subscription._acknowledgment;
        pendingSubscriptionHttpClient = null;
        return subscription;
      } catch (_) {
        pendingSubscriptionHttpClient = null;
        await subscription.close();
        rethrow;
      }
    },
    onAbort: abortSubscriptionSetup,
    trackForClose: false,
  );
}