listen method
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,
);
}