Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions packages/stream_core/CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,5 +1,9 @@
## Upcoming

### 💥 BREAKING CHANGES

- `AuthInterceptor` now takes a `TokenManager Function()` getter instead of a `TokenManager` instance. This lets callers swap the active `TokenManager` at runtime — e.g. after a guest token exchange resolves a server-assigned user id — and have the interceptor pick up the new instance (and its `userId`) on the next request.

### ✨ Features

- Added `teams` field to `User` class.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4,16 +4,31 @@ import '../../errors.dart';
import '../../user.dart';
import '../stream_core_dio_error.dart';

/// Provides the [TokenManager] currently in use by an [AuthInterceptor].
///
/// A getter rather than a fixed reference so the caller can swap the underlying
/// [TokenManager] at runtime — e.g. after a guest token exchange resolves a
/// server-assigned user id — and have the interceptor pick up the new instance.
typedef TokenManagerProvider = TokenManager Function();

/// Authentication interceptor that refreshes the token if
/// an auth error is received
class AuthInterceptor extends QueuedInterceptor {
/// Initialize a new auth interceptor
AuthInterceptor(this._dio, this._tokenManager);
/// Initialize a new auth interceptor.
///
/// [tokenManagerProvider] is a getter rather than a fixed reference so the
/// caller can swap the underlying [TokenManager] — e.g. after a guest token
/// exchange resolves a server-assigned user id — and have this interceptor
/// pick up the new instance on its next request.
AuthInterceptor(this._dio, this._tokenManagerProvider);

final Dio _dio;

/// The token manager used in the client
final TokenManager _tokenManager;
/// Provides the token manager currently in use.
final TokenManagerProvider _tokenManagerProvider;

/// The token manager currently in use.
TokenManager get _tokenManager => _tokenManagerProvider();

@override
Future<void> onRequest(
Expand All @@ -23,6 +38,11 @@ class AuthInterceptor extends QueuedInterceptor {
try {
final token = await _tokenManager.getToken();

// Re-read the token manager after awaiting the token: loading it may
// have swapped in a new manager carrying a server-resolved user id
// (e.g. a guest exchange). Reading `userId` here keeps the `user_id`
// query parameter consistent with the identity in the `Authorization`
// header below.
options.queryParameters['user_id'] = _tokenManager.userId;
options.headers['Authorization'] = token.rawValue;
options.headers['stream-auth-type'] = token.authType.headerValue;
Expand Down Expand Up @@ -57,10 +77,11 @@ class AuthInterceptor extends QueuedInterceptor {

final error = StreamApiError.fromJson(data);
if (error.isTokenExpiredError) {
final tokenManager = _tokenManager;
// Don't try to refresh the token if we're using a static provider
if (_tokenManager.usesStaticProvider) return handler.next(err);
if (tokenManager.usesStaticProvider) return handler.next(err);
// Otherwise, mark the current token as expired.
_tokenManager.expireToken();
tokenManager.expireToken();
Comment thread
coderabbitai[bot] marked this conversation as resolved.

try {
final options = err.requestOptions;
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,210 @@
import 'dart:convert';

import 'package:stream_core/stream_core.dart';
import 'package:test/test.dart';

// A minimal HttpClientAdapter that captures the outgoing RequestOptions and
// always responds with an empty successful response.
class _CapturingHttpClientAdapter implements HttpClientAdapter {
RequestOptions? lastRequest;

@override
Future<ResponseBody> fetch(
RequestOptions options,
Stream<Uint8List>? requestStream,
Future<void>? cancelFuture,
) async {
lastRequest = options;
return ResponseBody.fromString(
'{}',
200,
headers: {
Headers.contentTypeHeader: [Headers.jsonContentType],
},
);
}

@override
void close({bool force = false}) {}
}

// An adapter that always responds with a token-expired API error (code 40),
// counting how many times it is hit so a retry can be detected. [onFetch], if
// provided, runs when the request is dispatched — used to simulate a token
// manager being swapped in mid-flight.
class _TokenExpiredHttpClientAdapter implements HttpClientAdapter {
_TokenExpiredHttpClientAdapter({this.onFetch});

final void Function()? onFetch;

var _requestCount = 0;
int get requestCount => _requestCount;

@override
Future<ResponseBody> fetch(
RequestOptions options,
Stream<Uint8List>? requestStream,
Future<void>? cancelFuture,
) async {
_requestCount++;
onFetch?.call();
return ResponseBody.fromString(
jsonEncode({
'code': 40, // token expired
'details': <int>[],
'duration': '0ms',
'message': 'token expired',
'more_info': '',
'StatusCode': 401,
}),
401,
headers: {
Headers.contentTypeHeader: [Headers.jsonContentType],
},
);
}

@override
void close({bool force = false}) {}
}

UserToken _generateTestUserToken(String userId) {
String b64UrlNoPad(Object jsonObj) {
final bytes = utf8.encode(jsonEncode(jsonObj));
return base64Url.encode(bytes).replaceAll('=', '');
}

final header = {'alg': 'none', 'typ': 'JWT'};
final payload = {'user_id': userId};

final jwt = '${b64UrlNoPad(header)}.${b64UrlNoPad(payload)}.';
return UserToken(jwt);
}

void main() {
group('AuthInterceptor', () {
test(
'picks up a TokenManager swapped in while the token is loading, so the '
'user_id query parameter reflects a server-resolved id (guest exchange)',
() async {
// Simulates the guest flow: the token provider resolves to a
// server-assigned id and swaps in a new TokenManager carrying that id
// before the request headers are written. The interceptor reads the
// manager through the getter, so it observes the swapped instance.
late TokenManager tokenManager;
tokenManager = TokenManager(
userId: 'requested-id',
tokenProvider: TokenProvider.dynamic((_) async {
final token = _generateTestUserToken('server-assigned-id');
tokenManager = TokenManager(
userId: token.userId,
tokenProvider: TokenProvider.static(token),
);
return token;
}),
);

final dio = Dio(BaseOptions(baseUrl: 'https://example.com'));
final adapter = _CapturingHttpClientAdapter();
dio.httpClientAdapter = adapter;
dio.interceptors.add(AuthInterceptor(dio, () => tokenManager));

await dio.get<void>('/test');

expect(
adapter.lastRequest?.queryParameters['user_id'],
'server-assigned-id',
);
},
);

test(
'uses the current TokenManager userId when nothing swaps it '
'(regular/anonymous users)',
() async {
final tokenManager = TokenManager(
userId: 'user-123',
tokenProvider: TokenProvider.static(
_generateTestUserToken('user-123'),
),
);

final dio = Dio(BaseOptions(baseUrl: 'https://example.com'));
final adapter = _CapturingHttpClientAdapter();
dio.httpClientAdapter = adapter;
dio.interceptors.add(AuthInterceptor(dio, () => tokenManager));

await dio.get<void>('/test');

expect(adapter.lastRequest?.queryParameters['user_id'], 'user-123');
},
);

test(
'does not retry a token-expired response when using a static provider '
'(e.g. a guest token): the error is surfaced to the caller instead of '
'silently re-minting the token',
() async {
final tokenManager = TokenManager(
userId: 'guest-1',
tokenProvider: TokenProvider.static(_generateTestUserToken('guest-1')),
);

final dio = Dio(BaseOptions(baseUrl: 'https://example.com'));
final adapter = _TokenExpiredHttpClientAdapter();
dio.httpClientAdapter = adapter;
dio.interceptors.add(AuthInterceptor(dio, () => tokenManager));

await expectLater(
dio.get<void>('/test'),
throwsA(isA<DioException>()),
);

// A static provider must not trigger the refresh-and-retry path, so
// the request is attempted exactly once.
expect(adapter.requestCount, 1);
},
);

test(
'forwards a token-expired error without retrying when the token manager '
'is swapped to a static provider after the request was dispatched '
'(guest exchange resolving mid-flight)',
() async {
// Starts on a dynamic manager and swaps to a static one carrying the
// server-resolved id once the request is already in flight, mirroring
// the guest flow. onError observes the swapped-in (static) manager and
// must forward the error rather than expire + retry.
var tokenManager = TokenManager(
userId: 'requested-id',
tokenProvider: TokenProvider.dynamic(
(_) async => _generateTestUserToken('requested-id'),
),
);

final dio = Dio(BaseOptions(baseUrl: 'https://example.com'));
final adapter = _TokenExpiredHttpClientAdapter(
onFetch: () {
tokenManager = TokenManager(
userId: 'server-assigned-id',
tokenProvider: TokenProvider.static(
_generateTestUserToken('server-assigned-id'),
),
);
},
);
dio.httpClientAdapter = adapter;
dio.interceptors.add(AuthInterceptor(dio, () => tokenManager));

await expectLater(
dio.get<void>('/test'),
throwsA(isA<DioException>()),
);

// The swapped-in manager is static, so the error is surfaced without a
// refresh-and-retry: the request is attempted exactly once.
expect(adapter.requestCount, 1);
},
);
});
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,44 @@
import 'package:stream_core/stream_core.dart';
import 'package:test/test.dart';

StreamApiError _apiError(int code) => StreamApiError(
code: code,
details: const [],
duration: '0ms',
message: 'error $code',
moreInfo: '',
statusCode: 401,
);

Disconnected _serverDisconnect(StreamApiError apiError) => Disconnected(
source: ServerInitiated(
error: WebSocketEngineException(
reason: apiError.message,
code: 4001,
error: apiError,
),
),
);

void main() {
group('WebSocketConnectionState.isAutomaticReconnectionEnabled', () {
test(
'is disabled when the server closes with a token-expired error, so an '
'expired (e.g. guest) token does not trigger a silent reconnect loop',
() {
// Token-invalid error codes are 40..42; 40 = token expired.
final state = _serverDisconnect(_apiError(40));

expect(state.isAutomaticReconnectionEnabled, isFalse);
},
);

test('is enabled for a generic, retryable server-initiated disconnection', () {
// A server error that is neither a normal closure (1000), a token error
// (40..42), nor a client error (400..499) should still reconnect.
final state = _serverDisconnect(_apiError(43));

expect(state.isAutomaticReconnectionEnabled, isTrue);
});
});
}
Loading