diff --git a/README.md b/README.md index c99c19d..ee8f12f 100644 --- a/README.md +++ b/README.md @@ -10,6 +10,59 @@ This is the SDK for Flutter for [https://www.flagsmith.com/](https://www.flagsmi For full documentation visit [https://docs.flagsmith.com/clients/flutter/](https://docs.flagsmith.com/clients/flutter/) +## Experiments + +Once an experiment is running, Flagsmith serves the variations automatically through the flag. Your application +records exposures (when an identity experienced a variation) and conversion events (what your metrics aggregate). + +Enable event collection with `enableEvents`. Users must be identified: exposures and conversion events are joined per +identity, so use the same identifier for flags and events. + +```dart +final flagsmith = await FlagsmithClient.init( + apiKey: 'YOUR_CLIENT_SIDE_ENVIRONMENT_KEY', + config: const FlagsmithConfig(enableEvents: true), +); +final user = Identity(identifier: 'user_42'); +await flagsmith.getFeatureFlags(user: user); + +// Evaluate the flag and record an exposure in one call +final flag = await flagsmith.getExperimentFlag('checkout_button', user: user); +// ...render based on flag?.variant / flag?.stateValue + +// Record a conversion event; the name must match your metric's event name +flagsmith.trackEvent('purchase', value: 99.5); +``` + +### Exposures + +`getExperimentFlag` evaluates the flag and records a `$flag_exposure` event with the served `variant` as its value. It +is only recorded when the flag exists, is enabled and the identity is enrolled (`flag.experiment?.inExperiment == true`). +Anything else is logged and skipped, so it is safe to call against environments or servers without experiments. + +If you evaluate in one place and render in another, call `trackExposureEvent` at the point of display instead: + +```dart +final experiment = flag?.experiment; +if (flag != null && experiment != null) { + flagsmith.trackExposureEvent('checkout_button', + value: flag.variant, metadata: {'experiment_id': experiment.id}); +} +``` + +Exposures are deduplicated per identity and variant within a flush window, so recording one more than once is safe. + +### Conversion events + +`trackEvent` sends a named event, optionally with a `value`, `traits` and `metadata`. Names starting with `$` are +reserved. The event name is case-sensitive and must match the metric's configured event name. + +### Flushing + +Events are buffered and posted to `eventsURI` every `eventsFlushInterval` ms (10 s) or when `eventsMaxBuffer` (1000) +events are queued. `close()` flushes best-effort; await `flushEvents()` when you need the POST to complete, e.g. before +the app is torn down. Network failures are retried once, then logged and dropped; they never throw. + ## Contributing Please read [CONTRIBUTING.md](https://gist.github.com/kyle-ssg/c36a03aebe492e45cbd3eefb21cb0486) for details on our code of conduct, and the process for submitting pull requests to us. diff --git a/lib/src/core/core.dart b/lib/src/core/core.dart index f9f43a5..c116b0e 100644 --- a/lib/src/core/core.dart +++ b/lib/src/core/core.dart @@ -1,5 +1,6 @@ export 'crud_storage.dart'; export 'datetime_x.dart'; +export 'events/event_processor.dart'; export 'exceptions.dart'; export 'extensions/converters.dart'; export 'model/index.dart'; diff --git a/lib/src/core/events/event_processor.dart b/lib/src/core/events/event_processor.dart new file mode 100644 index 0000000..234792d --- /dev/null +++ b/lib/src/core/events/event_processor.dart @@ -0,0 +1,188 @@ +import 'dart:async'; +import 'dart:convert'; + +import 'package:dio/dio.dart'; + +import '../../version.dart'; + +/// Batches experimentation events to `{eventsURI}v1/events`. +/// +/// Flushes on a timer, at [maxBuffer], on [flush] and best-effort on [stop]. +/// Exposures dedupe within a flush window; a failed POST retries once then drops. +class EventProcessor { + static const String flagExposureEvent = r'$flag_exposure'; + static const String eventsPath = 'v1/events'; + static const String sdkUserAgentHeader = 'Flagsmith-SDK-User-Agent'; + static const String environmentKeyHeader = 'X-Environment-Key'; + static const String contentType = 'application/json; charset=utf-8'; + + final Dio _api; + final String _apiKey; + final String endpoint; + final int flushInterval; + final int maxBuffer; + final int retryBackoff; + final void Function(String message) _log; + + final List> _buffer = []; + final Set _dedupeKeys = {}; + final Set> _inFlight = {}; + Timer? _timer; + + EventProcessor({ + required Dio api, + required String apiKey, + required String eventsURI, + this.flushInterval = 10000, + this.maxBuffer = 1000, + this.retryBackoff = 1000, + void Function(String message)? log, + }) : _api = api, + _apiKey = apiKey, + _log = log ?? _noopLog, + endpoint = + '${eventsURI.endsWith('/') ? eventsURI : '$eventsURI/'}$eventsPath'; + + static void _noopLog(String _) {} + + List> get buffer => List.unmodifiable(_buffer); + + void trackEvent({ + required String event, + String? identifier, + Object? value, + Map? traits, + Map? metadata, + }) { + _bufferEvent( + event: event, + featureName: null, + identifier: identifier, + value: value, + traits: traits, + metadata: metadata, + dedupe: false, + ); + } + + void trackExposureEvent({ + required String featureName, + required String identifier, + Object? value, + Map? traits, + Map? metadata, + }) { + _bufferEvent( + event: flagExposureEvent, + featureName: featureName, + identifier: identifier, + value: value, + traits: traits, + metadata: metadata, + dedupe: true, + ); + } + + void _bufferEvent({ + required String event, + required String? featureName, + required String? identifier, + required Object? value, + required Map? traits, + required Map? metadata, + required bool dedupe, + }) { + final stringValue = value == null ? null : '$value'; + if (dedupe) { + // Experiment id is part of the key so a new experiment on the same flag + // and variant within one flush window still records its own exposure. + final key = jsonEncode([ + event, + featureName, + identifier, + stringValue, + metadata?['experiment_id'] + ]); + if (_dedupeKeys.contains(key)) { + return; + } + _dedupeKeys.add(key); + } + _buffer.add({ + 'event': event, + 'feature_name': featureName, + 'identifier': identifier, + 'value': stringValue, + 'traits': traits, + 'metadata': { + ...?metadata, + 'sdk_version': sdkVersion, + }, + 'timestamp': DateTime.now().millisecondsSinceEpoch, + }); + if (_buffer.length >= maxBuffer) { + unawaited(flush()); + } + } + + /// Posts the buffered events and waits for every upload still in flight, + /// including ones started by the timer or the max-buffer trigger, so that + /// awaiting it at teardown means the POSTs have completed. Never throws. + Future flush() async { + if (_buffer.isNotEmpty) { + final events = List>.from(_buffer); + _buffer.clear(); + _dedupeKeys.clear(); + late final Future upload; + upload = + _postBatch(events, 0).whenComplete(() => _inFlight.remove(upload)); + _inFlight.add(upload); + } + await Future.wait(_inFlight.toList()); + } + + void start() { + _timer?.cancel(); + _timer = null; + if (flushInterval > 0) { + _timer = Timer.periodic( + Duration(milliseconds: flushInterval), (_) => unawaited(flush())); + } + } + + /// Cancels the timer and flushes without awaiting; await [flush] for teardown. + void stop() { + _timer?.cancel(); + _timer = null; + unawaited(flush()); + } + + Future _postBatch( + List> events, int attempt) async { + try { + final response = await _api.post( + endpoint, + data: {'events': events}, + options: Options( + contentType: contentType, + headers: { + environmentKeyHeader: _apiKey, + sdkUserAgentHeader: getUserAgent(), + }, + ), + ); + final status = response.statusCode ?? 0; + if (status < 200 || status >= 300) { + throw StateError('unexpected status $status'); + } + _log('Events: flush successful (${events.length} events)'); + } catch (e) { + if (attempt < 1) { + _log('Events: flush failed, retrying: $e'); + await Future.delayed(Duration(milliseconds: retryBackoff)); + return _postBatch(events, attempt + 1); + } + _log('Events: flush failed, dropping ${events.length} events: $e'); + } + } +} diff --git a/lib/src/core/model/experiment.dart b/lib/src/core/model/experiment.dart new file mode 100644 index 0000000..95281c8 --- /dev/null +++ b/lib/src/core/model/experiment.dart @@ -0,0 +1,51 @@ +import 'package:json_annotation/json_annotation.dart'; + +part 'experiment.g.dart'; + +/// The running experiment a flag was evaluated under; identity evaluations only. +@JsonSerializable() +class Experiment { + final int id; + final String name; + + /// Whether the identity is enrolled. `variant` alone cannot tell. + @JsonKey(name: 'in_experiment', defaultValue: false) + final bool inExperiment; + + const Experiment({ + required this.id, + required this.name, + this.inExperiment = false, + }); + + factory Experiment.fromJson(Map json) => + _$ExperimentFromJson(json); + + Map toJson() => _$ExperimentToJson(this); + + @override + String toString() => 'Experiment($id:$name, inExperiment=$inExperiment)'; +} + +/// `metadata.experiment` -> [Experiment]; null when absent or malformed. +Experiment? experimentFromMetadata(Object? metadata) { + if (metadata is! Map) { + return null; + } + final experiment = metadata['experiment']; + if (experiment is! Map) { + return null; + } + try { + return Experiment.fromJson(Map.from(experiment)); + } catch (_) { + return null; + } +} + +Map? experimentToMetadata(Experiment? experiment) { + if (experiment == null) { + return null; + } + return {'experiment': experiment.toJson()}; +} diff --git a/lib/src/core/model/experiment.g.dart b/lib/src/core/model/experiment.g.dart new file mode 100644 index 0000000..a0dc843 --- /dev/null +++ b/lib/src/core/model/experiment.g.dart @@ -0,0 +1,22 @@ +// GENERATED CODE - DO NOT MODIFY BY HAND + +// ignore_for_file: implicit_dynamic_parameter, non_constant_identifier_names, type_annotate_public_apis, omit_local_variable_types, unnecessary_this + +part of 'experiment.dart'; + +// ************************************************************************** +// JsonSerializableGenerator +// ************************************************************************** + +Experiment _$ExperimentFromJson(Map json) => Experiment( + id: (json['id'] as num).toInt(), + name: json['name'] as String, + inExperiment: json['in_experiment'] as bool? ?? false, + ); + +Map _$ExperimentToJson(Experiment instance) => + { + 'id': instance.id, + 'name': instance.name, + 'in_experiment': instance.inExperiment, + }; diff --git a/lib/src/core/model/flag.dart b/lib/src/core/model/flag.dart index 9f15799..8e2d0ca 100644 --- a/lib/src/core/model/flag.dart +++ b/lib/src/core/model/flag.dart @@ -4,6 +4,7 @@ import '../extensions/converters.dart'; import 'package:json_annotation/json_annotation.dart'; import 'dart:math'; +import 'experiment.dart'; import 'feature.dart'; part 'flag.g.dart'; @@ -21,6 +22,17 @@ class Flag { final int? identity; @JsonKey(name: 'feature_segment') final int? featureSegment; + + final String? variant; + final String? reason; + + /// Lifted from `metadata.experiment`; null unless an experiment is running. + @JsonKey( + name: 'metadata', + fromJson: experimentFromMetadata, + toJson: experimentToMetadata, + includeIfNull: false) + final Experiment? experiment; Flag( {this.id, required this.feature, @@ -28,7 +40,10 @@ class Flag { this.enabled, this.environment, this.identity, - this.featureSegment}); + this.featureSegment, + this.variant, + this.reason, + this.experiment}); String get key => feature.name; @override @@ -45,7 +60,10 @@ class Flag { bool? enabled, int? environment, int? identity, - int? featureSegment}) => + int? featureSegment, + String? variant, + String? reason, + Experiment? experiment}) => Flag( id: id, feature: feature, @@ -54,6 +72,9 @@ class Flag { environment: environment, identity: identity, featureSegment: featureSegment, + variant: variant, + reason: reason, + experiment: experiment, ); factory Flag.seed(String featureName, {bool enabled = true, String? value}) { var id = _generateNum(1, 100); @@ -85,7 +106,10 @@ class Flag { bool? enabled, int? environment, int? identity, - int? featureSegment}) => + int? featureSegment, + String? variant, + String? reason, + Experiment? experiment}) => Flag( id: id ?? this.id, feature: feature ?? this.feature, @@ -94,5 +118,8 @@ class Flag { environment: environment ?? this.environment, identity: identity ?? this.identity, featureSegment: featureSegment ?? this.featureSegment, + variant: variant ?? this.variant, + reason: reason ?? this.reason, + experiment: experiment ?? this.experiment, ); } diff --git a/lib/src/core/model/flag.g.dart b/lib/src/core/model/flag.g.dart index ddfecd9..d0d7bee 100644 --- a/lib/src/core/model/flag.g.dart +++ b/lib/src/core/model/flag.g.dart @@ -16,14 +16,30 @@ Flag _$FlagFromJson(Map json) => Flag( environment: (json['environment'] as num?)?.toInt(), identity: (json['identity'] as num?)?.toInt(), featureSegment: (json['feature_segment'] as num?)?.toInt(), + variant: json['variant'] as String?, + reason: json['reason'] as String?, + experiment: experimentFromMetadata(json['metadata']), ); -Map _$FlagToJson(Flag instance) => { - 'id': instance.id, - 'feature': instance.feature.toJson(), - 'feature_state_value': stringToJson(instance.stateValue), - 'enabled': instance.enabled, - 'environment': instance.environment, - 'identity': instance.identity, - 'feature_segment': instance.featureSegment, - }; +Map _$FlagToJson(Flag instance) { + final val = { + 'id': instance.id, + 'feature': instance.feature.toJson(), + 'feature_state_value': stringToJson(instance.stateValue), + 'enabled': instance.enabled, + 'environment': instance.environment, + 'identity': instance.identity, + 'feature_segment': instance.featureSegment, + 'variant': instance.variant, + 'reason': instance.reason, + }; + + void writeNotNull(String key, dynamic value) { + if (value != null) { + val[key] = value; + } + } + + writeNotNull('metadata', experimentToMetadata(instance.experiment)); + return val; +} diff --git a/lib/src/core/model/index.dart b/lib/src/core/model/index.dart index c9aaef4..135353e 100644 --- a/lib/src/core/model/index.dart +++ b/lib/src/core/model/index.dart @@ -1,5 +1,6 @@ library; +export 'experiment.dart'; export 'feature.dart'; export 'identity.dart'; export 'flag.dart'; diff --git a/lib/src/flagsmith_client.dart b/lib/src/flagsmith_client.dart index 3a2c4bb..1bd97f9 100644 --- a/lib/src/flagsmith_client.dart +++ b/lib/src/flagsmith_client.dart @@ -46,6 +46,9 @@ class FlagsmithClient { final Map flagAnalytics = {}; Timer? _analyticsTimer; + EventProcessor? _eventProcessor; + EventProcessor? get eventProcessor => _eventProcessor; + final StreamController _loading = StreamController.broadcast(); @@ -67,6 +70,16 @@ class FlagsmithClient { if (config.enableAnalytics) { _setupAnalyticsTimer(config.analyticsInterval); } + if (config.enableEvents) { + _eventProcessor = EventProcessor( + api: _api, + apiKey: apiKey, + eventsURI: config.eventsURI, + flushInterval: config.eventsFlushInterval, + maxBuffer: config.eventsMaxBuffer, + log: log, + )..start(); + } if (config.enableRealtimeUpdates) { _setupRealtimeUpdates(config.realtimeUpdatesBaseURI); } @@ -169,7 +182,8 @@ class FlagsmithClient { switch (config.storageType) { case StorageType.custom: if (storage == null) { - throw FlagsmithConfigException(Exception('When using StorageType.custom, a storage implementation must be provided')); + throw FlagsmithConfigException(Exception( + 'When using StorageType.custom, a storage implementation must be provided')); } store = storage; break; @@ -346,6 +360,114 @@ class FlagsmithClient { return feature?.stateValue; } + /// EXPERIMENTS + /// + /// Resolve a flag for [user] (or [cachedUser]) and fire one `$flag_exposure` + /// event with the variant as value. Skipped unless events are enabled, the + /// flag is enabled and `flag.experiment.inExperiment` is true. + /// + /// When [user] is supplied the flags are fetched for that identity unless + /// [reload] is explicitly false, so the exposure never reuses another + /// identity's stored assignment. Without [user], stored flags are used. + Future getExperimentFlag(String featureName, + {Identity? user, List? traits, bool? reload}) async { + final identity = user ?? cachedUser; + if (identity != null) { + cachedUser = identity; + } + final flags = await getFeatureFlags( + user: identity, traits: traits, reload: reload ?? (user != null)); + final flag = flags + .firstWhereOrNull((element) => element.feature.name == featureName); + _incrementFlagAnalytics(flag); + + if (_eventProcessor == null) { + return flag; + } + if (identity == null) { + log('getExperimentFlag called for "$featureName" without an identity. ' + 'No exposure recorded.'); + return flag; + } + if (flag == null) { + log('getExperimentFlag called for "$featureName" which does not exist. ' + 'No exposure recorded.'); + return null; + } + if (flag.enabled != true) { + log('getExperimentFlag called for "$featureName" which is disabled. ' + 'No exposure recorded.'); + return flag; + } + final experiment = flag.experiment; + if (experiment == null || !experiment.inExperiment) { + log('getExperimentFlag called for "$featureName" but this identity is ' + 'not enrolled in a running experiment for it. No exposure recorded.'); + return flag; + } + trackExposureEvent(featureName, + user: identity, + value: flag.variant, + metadata: {'experiment_id': experiment.id}); + return flag; + } + + /// Record a conversion event. No-op unless events are enabled; names + /// starting with `$` are reserved and throw [ArgumentError]. + void trackEvent(String event, + {Identity? user, + Object? value, + Map? traits, + Map? metadata}) { + if (event.startsWith(r'$')) { + throw ArgumentError.value( + event, + 'event', + r'event names starting with "$" are reserved; ' + 'use trackExposureEvent to record an exposure'); + } + final processor = _eventProcessor; + if (processor == null) { + return; + } + processor.trackEvent( + event: event, + identifier: (user ?? cachedUser)?.identifier, + value: value, + traits: traits, + metadata: metadata, + ); + } + + /// Record a `$flag_exposure` event at the point of display. No-op unless + /// events are enabled; requires an identity ([user] or [cachedUser]). + void trackExposureEvent(String featureName, + {Identity? user, + Object? value, + Map? traits, + Map? metadata}) { + final processor = _eventProcessor; + if (processor == null) { + return; + } + final identifier = (user ?? cachedUser)?.identifier; + if (identifier == null) { + log('trackExposureEvent called for "$featureName" without an identity. ' + 'No exposure recorded.'); + return; + } + processor.trackExposureEvent( + featureName: featureName, + identifier: identifier, + value: value, + traits: traits, + metadata: metadata, + ); + } + + /// Flush buffered events and await the POST. Never throws. + Future flushEvents() => _eventProcessor?.flush() ?? Future.value(); + /// Internal function for collecting analytical data on flag usage void _incrementFlagAnalytics(Flag? flag) { if (flag != null && config.enableAnalytics) { @@ -569,6 +691,7 @@ class FlagsmithClient { void close() { _analyticsTimer?.cancel(); + _eventProcessor?.stop(); SSEClient.unsubscribeFromSSE(); _loading.close(); } diff --git a/lib/src/flagsmith_config.dart b/lib/src/flagsmith_config.dart index 7691068..851fa43 100644 --- a/lib/src/flagsmith_config.dart +++ b/lib/src/flagsmith_config.dart @@ -27,6 +27,11 @@ class FlagsmithConfig { final String realtimeUpdatesBaseURI; final int reconnectToSSEInterval; + final bool enableEvents; + final String eventsURI; + final int eventsFlushInterval; + final int eventsMaxBuffer; + /// Flagsmith config initialization /// change only if you have self-hosted Flagsmith /// [baseURI], [flagsURI], [identitiesURI], [traitsURI], [analyticsURI] @@ -51,6 +56,9 @@ class FlagsmithConfig { /// If you would like to use realtime updates, set [enableRealtimeUpdates] to *true* /// /// You can configure the realtime updates source URL by setting the [realtimeUpdatesBaseURI] parameter + /// + /// Set [enableEvents] to *true* to batch experiment events to [eventsURI], + /// flushed every [eventsFlushInterval] ms or at [eventsMaxBuffer] events const FlagsmithConfig({ this.baseURI = 'https://edge.api.flagsmith.com/api/v1/', @@ -71,6 +79,10 @@ class FlagsmithConfig { this.realtimeUpdatesBaseURI = 'https://realtime.flagsmith.com/sse/environments/', this.reconnectToSSEInterval = 29000, + this.enableEvents = false, + this.eventsURI = 'https://events.api.flagsmith.com/', + this.eventsFlushInterval = 10000, + this.eventsMaxBuffer = 1000, }); /// Client options from config diff --git a/test/core/models/flag_test.dart b/test/core/models/flag_test.dart index 65ed07d..113b4c0 100644 --- a/test/core/models/flag_test.dart +++ b/test/core/models/flag_test.dart @@ -93,6 +93,97 @@ void main() { }); }); + group('[Experiment]', () { + Map flagJson({Object? metadata, bool withKey = true}) => + { + 'id': 7, + 'feature': {'id': 7, 'name': 'checkout_cta'}, + 'enabled': true, + 'feature_state_value': 'buy-now', + 'variant': 'treatment-a', + 'reason': 'SPLIT; weight=50', + if (withKey) 'metadata': metadata, + }; + + test('When metadata.experiment present, then experiment is populated', () { + final flag = Flag.fromJson(flagJson(metadata: { + 'experiment': {'id': 42, 'name': 'New checkout CTA', 'in_experiment': true} + })); + expect(flag.variant, 'treatment-a'); + expect(flag.reason, 'SPLIT; weight=50'); + expect(flag.experiment, isNotNull); + expect(flag.experiment!.id, 42); + expect(flag.experiment!.name, 'New checkout CTA'); + expect(flag.experiment!.inExperiment, isTrue); + }); + + test('When in_experiment false, then inExperiment is false', () { + final flag = Flag.fromJson(flagJson(metadata: { + 'experiment': {'id': 42, 'name': 'x', 'in_experiment': false} + })); + expect(flag.experiment!.inExperiment, isFalse); + }); + + test('When metadata absent or null, then experiment is null', () { + expect(Flag.fromJson(flagJson(withKey: false)).experiment, isNull); + expect(Flag.fromJson(flagJson(metadata: null)).experiment, isNull); + }); + + test('When metadata has only unknown keys, then experiment is null', () { + final flag = Flag.fromJson(flagJson(metadata: { + 'something_else': {'id': 1} + })); + expect(flag.experiment, isNull); + }); + + test('When experiment is malformed, then parsing still succeeds', () { + final flag = Flag.fromJson(flagJson(metadata: { + 'experiment': {'name': 'missing id'} + })); + expect(flag.experiment, isNull); + expect(flag.variant, 'treatment-a'); + }); + + test('When old server omits variant and reason, then both are null', () { + final flag = Flag.fromJson({ + 'feature': {'id': 7, 'name': 'checkout_cta'}, + 'enabled': true, + 'feature_state_value': null, + }); + expect(flag.variant, isNull); + expect(flag.reason, isNull); + expect(flag.experiment, isNull); + expect(flag.toJson().containsKey('metadata'), isFalse); + }); + + test('When flag round-trips through toJson, then experiment survives', () { + final original = Flag.fromJson(flagJson(metadata: { + 'experiment': {'id': 42, 'name': 'New checkout CTA', 'in_experiment': true}, + 'unknown_key': 1, + })); + final json = original.toJson(); + expect(json['metadata'], { + 'experiment': {'id': 42, 'name': 'New checkout CTA', 'in_experiment': true} + }); + + final restored = Flag.fromJson(jsonDecode(jsonEncode(json))); + expect(restored.variant, original.variant); + expect(restored.reason, original.reason); + expect(restored.experiment!.id, 42); + expect(restored.experiment!.name, 'New checkout CTA'); + expect(restored.experiment!.inExperiment, isTrue); + }); + + test('When copyWith sets experiment, then other fields are kept', () { + final flag = Flag.fromJson(flagJson(withKey: false)); + final copy = flag.copyWith( + experiment: const Experiment(id: 1, name: 'e', inExperiment: true)); + expect(copy.experiment!.id, 1); + expect(copy.variant, 'treatment-a'); + expect(flag.experiment, isNull); + }); + }); + group('[FlagAndTraits]', () { test('When response successfuly parsed', () { final identity = FlagsAndTraits.fromJson( diff --git a/test/fg/flagsmith_experiments_test.dart b/test/fg/flagsmith_experiments_test.dart new file mode 100644 index 0000000..c70ce98 --- /dev/null +++ b/test/fg/flagsmith_experiments_test.dart @@ -0,0 +1,529 @@ +import 'dart:convert'; + +import 'package:dio/dio.dart'; +import 'package:flagsmith/flagsmith.dart'; +import 'package:flagsmith/src/version.dart'; +import 'package:http_mock_adapter/http_mock_adapter.dart'; +import 'package:test/test.dart'; + +import '../shared.dart'; + +const eventsEndpoint = 'https://events.api.flagsmith.com/v1/events'; +const user = Identity(identifier: 'user_42'); + +/// Records every request Dio sends to the events endpoint. +class EventsCapture { + final List requests = []; + + List> get events => requests + .expand((r) => (r.data['events'] as List).cast>()) + .toList(); + + Interceptor get interceptor => InterceptorsWrapper(onRequest: (options, h) { + if (options.path == eventsEndpoint) { + requests.add(options); + } + h.next(options); + }); +} + +Future<(FlagsmithClient, EventsCapture)> buildClient({ + bool enableEvents = true, + int flushInterval = 60000, + int maxBuffer = 1000, + bool eventsFail = false, + bool loadUserFlags = true, +}) async { + final fs = FlagsmithClient( + apiKey: apiKey, + seeds: seeds, + config: FlagsmithConfig( + baseURI: 'https://offline.net/', + enableEvents: enableEvents, + eventsFlushInterval: flushInterval, + eventsMaxBuffer: maxBuffer, + ), + ); + final capture = EventsCapture(); + fs.client.interceptors.add(capture.interceptor); + setupAdapter(fs, cb: (config, adapter) { + adapter.onPost(config.identitiesURI, (server) { + server.reply(200, jsonDecode(identitiesResponseData)); + }, data: Matchers.any); + adapter.onPost(eventsEndpoint, (server) { + if (eventsFail) { + server.throws(500, + DioException(requestOptions: RequestOptions(path: eventsEndpoint))); + return; + } + server.reply(200, {}); + }, data: Matchers.any); + }); + await fs.initialize(); + if (loadUserFlags) { + await fs.getFeatureFlags(user: user); + } + return (fs, capture); +} + +void main() { + group('[Experiments] getExperimentFlag', () { + late FlagsmithClient fs; + late EventsCapture capture; + setUp(() async { + (fs, capture) = await buildClient(); + }); + tearDown(() => fs.close()); + + test('When identity is enrolled, then flag is returned and exposure posted', + () async { + final flag = await fs.getExperimentFlag(experimentFeatureName); + + expect(flag, isNotNull); + expect(flag!.enabled, isTrue); + expect(flag.stateValue, 'buy-now'); + expect(flag.variant, experimentVariant); + expect(flag.reason, 'SPLIT; weight=50'); + expect(flag.experiment!.id, experimentId); + expect(flag.experiment!.inExperiment, isTrue); + expect(fs.eventProcessor!.buffer.length, 1); + + await fs.flushEvents(); + + expect(capture.requests.length, 1); + final request = capture.requests.single; + expect(request.method, 'POST'); + expect(request.contentType, 'application/json; charset=utf-8'); + expect(request.headers['X-Environment-Key'], apiKey); + expect(request.headers['Flagsmith-SDK-User-Agent'], + 'flagsmith-flutter-sdk/$sdkVersion'); + + final event = capture.events.single; + expect(event['event'], r'$flag_exposure'); + expect(event['feature_name'], experimentFeatureName); + expect(event['identifier'], user.identifier); + expect(event['value'], experimentVariant); + expect(event['traits'], isNull); + expect(event['metadata'], { + 'experiment_id': experimentId, + 'sdk_version': sdkVersion, + }); + expect(event['timestamp'], isA()); + expect(fs.eventProcessor!.buffer, isEmpty); + }); + + test('When identity is not enrolled, then no exposure', () async { + final flag = await fs.getExperimentFlag(experimentNotEnrolledFeatureName); + expect(flag!.variant, 'control'); + expect(flag.experiment!.inExperiment, isFalse); + expect(fs.eventProcessor!.buffer, isEmpty); + }); + + test('When flag has no experiment metadata, then no exposure', () async { + final flag = await fs.getExperimentFlag(experimentNoMetadataFeatureName); + expect(flag!.variant, 'control'); + expect(flag.experiment, isNull); + expect(fs.eventProcessor!.buffer, isEmpty); + }); + + test('When flag is disabled, then flag returned and no exposure', () async { + final flag = await fs.getExperimentFlag(experimentDisabledFeatureName); + expect(flag!.enabled, isFalse); + expect(flag.experiment!.inExperiment, isTrue); + expect(fs.eventProcessor!.buffer, isEmpty); + }); + + test('When flag is missing, then null and no exposure', () async { + final flag = await fs.getExperimentFlag(notImplementedFeatureName); + expect(flag, isNull); + expect(fs.eventProcessor!.buffer, isEmpty); + }); + + test('When no identity is known, then flag returned and no exposure', + () async { + fs.cachedUser = null; + final flag = await fs.getExperimentFlag(experimentFeatureName); + expect(flag!.experiment!.inExperiment, isTrue); + expect(fs.eventProcessor!.buffer, isEmpty); + }); + + test('When user param is given, then its flags are fetched and it wins', + () async { + // other_user gets a different variant from the server than user_42. + final identityPosts = []; + fs.client.interceptors.add(InterceptorsWrapper(onRequest: (o, h) { + if (o.path == fs.config.identitiesURI) { + final id = (o.data as Map)['identifier'] as String; + identityPosts.add(id); + if (id == 'other_user') { + final body = jsonDecode(identitiesResponseData); + for (final f in body['flags'] as List) { + if (f['feature']['name'] == experimentFeatureName) { + f['variant'] = 'treatment-b'; + } + } + h.resolve(Response(requestOptions: o, statusCode: 200, data: body)); + return; + } + } + h.next(o); + })); + + final flag = await fs.getExperimentFlag(experimentFeatureName, + user: const Identity(identifier: 'other_user')); + + expect(identityPosts, ['other_user'], + reason: 'must fetch B, not reuse A'); + expect(flag!.variant, 'treatment-b'); + final event = fs.eventProcessor!.buffer.single; + expect(event['identifier'], 'other_user'); + expect(event['value'], 'treatment-b'); + expect(fs.cachedUser?.identifier, 'other_user'); + }); + + test('When user param is given with reload false, then storage is used', + () async { + final flag = await fs.getExperimentFlag(experimentFeatureName, + user: const Identity(identifier: 'other_user'), reload: false); + expect(flag!.variant, experimentVariant); + expect(fs.eventProcessor!.buffer.single['identifier'], 'other_user'); + }); + + test('When called twice in a window, then exposure is deduped', () async { + await fs.getExperimentFlag(experimentFeatureName); + await fs.getExperimentFlag(experimentFeatureName); + expect(fs.eventProcessor!.buffer.length, 1); + + await fs.flushEvents(); + await fs.getExperimentFlag(experimentFeatureName); + expect(fs.eventProcessor!.buffer.length, 1); + }); + + test('When flag is read, then flag analytics are incremented', () async { + await fs.getExperimentFlag(experimentFeatureName); + await fs.getExperimentFlag(experimentNotEnrolledFeatureName); + await fs.getExperimentFlag(experimentNotEnrolledFeatureName); + expect(fs.flagAnalytics[experimentFeatureName], 1); + expect(fs.flagAnalytics[experimentNotEnrolledFeatureName], 2); + }); + + test('When flag restored from storage, then it still gates', () async { + await fs.getExperimentFlag(experimentFeatureName); + await fs.flushEvents(); + final restored = await fs.storageProvider.read(experimentFeatureName); + expect(restored!.experiment!.inExperiment, isTrue); + expect(restored.variant, experimentVariant); + }); + }); + + group('[Experiments] events disabled', () { + late FlagsmithClient fs; + late EventsCapture capture; + setUp(() async { + (fs, capture) = await buildClient(enableEvents: false); + }); + tearDown(() => fs.close()); + + test('When events disabled, then plain read and nothing posted', () async { + expect(fs.config.enableEvents, isFalse); + expect(fs.eventProcessor, isNull); + + final flag = await fs.getExperimentFlag(experimentFeatureName); + expect(flag!.experiment!.inExperiment, isTrue); + + fs.trackEvent('purchase', value: 1); + fs.trackExposureEvent(experimentFeatureName, value: 'x'); + await fs.flushEvents(); + expect(capture.requests, isEmpty); + }); + + test('When events disabled, then reserved names still throw', () { + expect(() => fs.trackEvent(r'$flag_exposure'), throwsArgumentError); + }); + }); + + group('[Experiments] trackEvent and trackExposureEvent', () { + late FlagsmithClient fs; + late EventsCapture capture; + setUp(() async { + (fs, capture) = await buildClient(); + }); + tearDown(() => fs.close()); + + test('When event name starts with \$, then throws', () { + expect(() => fs.trackEvent(r'$x'), throwsArgumentError); + expect(fs.eventProcessor!.buffer, isEmpty); + }); + + test('When conversion tracked, then buffered with stringified value', + () async { + fs.trackEvent('purchase', value: 99.5); + final event = fs.eventProcessor!.buffer.single; + expect(event['event'], 'purchase'); + expect(event['feature_name'], isNull); + expect(event['identifier'], user.identifier); + expect(event['value'], '99.5'); + expect(event['metadata'], {'sdk_version': sdkVersion}); + + await fs.flushEvents(); + expect(capture.events.single['event'], 'purchase'); + }); + + test('When conversion tracked twice, then never deduped', () { + fs.trackEvent('purchase', value: 1); + fs.trackEvent('purchase', value: 1); + expect(fs.eventProcessor!.buffer.length, 2); + }); + + test('When traits, metadata and user given, then passed through', () { + fs.trackEvent('signup', + user: const Identity(identifier: 'u2'), + traits: {'plan': 'premium'}, + metadata: {'source': 'qa'}); + final event = fs.eventProcessor!.buffer.single; + expect(event['identifier'], 'u2'); + expect(event['traits'], {'plan': 'premium'}); + expect(event['metadata'], {'source': 'qa', 'sdk_version': sdkVersion}); + expect(event['value'], isNull); + }); + + test('When exposure tracked without identity, then dropped', () { + fs.cachedUser = null; + fs.trackExposureEvent(experimentFeatureName, value: 'treatment-a'); + expect(fs.eventProcessor!.buffer, isEmpty); + }); + + test('When exposure tracked manually, then deduped per value', () { + fs.trackExposureEvent(experimentFeatureName, value: 'treatment-a'); + fs.trackExposureEvent(experimentFeatureName, value: 'treatment-a'); + fs.trackExposureEvent(experimentFeatureName, value: 'treatment-b'); + fs.trackExposureEvent(experimentFeatureName, + value: 'treatment-a', user: const Identity(identifier: 'u2')); + expect(fs.eventProcessor!.buffer.length, 3); + expect(fs.eventProcessor!.buffer.first['event'], r'$flag_exposure'); + }); + + test( + 'When same variant is exposed under a new experiment, then not deduped', + () { + fs.trackExposureEvent(experimentFeatureName, + value: 'treatment-a', metadata: {'experiment_id': 42}); + fs.trackExposureEvent(experimentFeatureName, + value: 'treatment-a', metadata: {'experiment_id': 42}); + fs.trackExposureEvent(experimentFeatureName, + value: 'treatment-a', metadata: {'experiment_id': 43}); + fs.trackExposureEvent(experimentFeatureName, value: 'treatment-a'); + expect( + fs.eventProcessor!.buffer.map((e) => e['metadata']['experiment_id']), + [42, 43, null]); + }); + }); + + group('[Experiments] flush', () { + test('When flush interval elapses, then buffer is posted', () async { + final (fs, capture) = await buildClient(flushInterval: 50); + fs.trackEvent('purchase'); + expect(capture.requests, isEmpty); + + await Future.delayed(const Duration(milliseconds: 250)); + expect(capture.requests.length, 1); + expect(fs.eventProcessor!.buffer, isEmpty); + fs.close(); + }); + + test('When max buffer is reached, then buffer is posted', () async { + final (fs, capture) = await buildClient(maxBuffer: 2); + fs.trackEvent('one'); + expect(capture.requests, isEmpty); + fs.trackEvent('two'); + + await Future.delayed(const Duration(milliseconds: 50)); + expect(capture.events.map((e) => e['event']), ['one', 'two']); + expect(fs.eventProcessor!.buffer, isEmpty); + fs.close(); + }); + + test('When flushEvents called on empty buffer, then nothing posted', + () async { + final (fs, capture) = await buildClient(); + await fs.flushEvents(); + expect(capture.requests, isEmpty); + fs.close(); + }); + + test('When client is closed, then buffer is flushed', () async { + final (fs, capture) = await buildClient(); + fs.trackEvent('purchase'); + fs.close(); + + await Future.delayed(const Duration(milliseconds: 50)); + expect(capture.events.single['event'], 'purchase'); + }); + + test('When POST fails, then flushEvents does not throw and drops batch', + () async { + final (fs, capture) = await buildClient(eventsFail: true); + fs.trackEvent('purchase'); + + await expectLater(fs.flushEvents(), completes); + expect(capture.requests.length, 2, reason: 'one retry then drop'); + expect(fs.eventProcessor!.buffer, isEmpty); + fs.close(); + }); + }); + + group('[Experiments] EventProcessor', () { + late Dio dio; + late DioAdapter adapter; + late EventsCapture capture; + late List logs; + + setUp(() { + dio = Dio(BaseOptions(baseUrl: 'https://offline.net/')); + adapter = DioAdapter(dio: dio); + capture = EventsCapture(); + dio.interceptors.add(capture.interceptor); + logs = []; + }); + + EventProcessor processor({int retryBackoff = 10, int flushInterval = 0}) => + EventProcessor( + api: dio, + apiKey: apiKey, + eventsURI: 'https://events.api.flagsmith.com', + retryBackoff: retryBackoff, + flushInterval: flushInterval, + log: logs.add, + ); + + test('When eventsURI lacks trailing slash, then endpoint is normalised', + () { + expect(processor().endpoint, eventsEndpoint); + }); + + test('When POST returns non-2xx, then retried once and dropped', () async { + adapter.onPost(eventsEndpoint, (server) { + server.reply(500, {}); + }, data: Matchers.any); + + final p = processor(); + p.trackEvent(event: 'purchase', identifier: 'u1'); + await p.flush(); + + expect(capture.requests.length, 2); + expect(p.buffer, isEmpty); + expect(logs.where((l) => l.contains('retrying')).length, 1); + expect(logs.where((l) => l.contains('dropping')).length, 1); + }); + + test('When POST fails once then succeeds, then batch is delivered', + () async { + adapter.onPost(eventsEndpoint, (server) { + server.throws( + 0, + DioException.connectionError( + requestOptions: RequestOptions(path: eventsEndpoint), + reason: 'offline')); + }, data: Matchers.any); + + final p = processor(retryBackoff: 100); + p.trackEvent(event: 'purchase', identifier: 'u1'); + final flushing = p.flush(); + + // Last registered handler wins: swap to success during the backoff. + await Future.delayed(const Duration(milliseconds: 20)); + adapter.onPost(eventsEndpoint, (server) { + server.reply(200, {}); + }, data: Matchers.any); + await flushing; + + expect(capture.requests.length, 2); + expect(capture.events.map((e) => e['event']), ['purchase', 'purchase']); + expect(logs.last, contains('flush successful')); + }); + + test('When events tracked during a flush, then they wait for the next one', + () async { + adapter.onPost(eventsEndpoint, (server) { + server.reply(200, {}, + delay: const Duration(milliseconds: 30)); + }, data: Matchers.any); + + final p = processor(); + p.trackEvent(event: 'first'); + final flushing = p.flush(); + p.trackEvent(event: 'second'); + await flushing; + + expect(capture.events.map((e) => e['event']), ['first']); + expect(p.buffer.single['event'], 'second'); + }); + + test('When max-buffer flush is in flight, then flush awaits its POST', + () async { + var responses = 0; + dio.interceptors.add(InterceptorsWrapper(onResponse: (r, h) { + responses++; + h.next(r); + })); + adapter.onPost(eventsEndpoint, (server) { + server.reply(200, {}, + delay: const Duration(milliseconds: 80)); + }, data: Matchers.any); + + final p = EventProcessor( + api: dio, + apiKey: apiKey, + eventsURI: 'https://events.api.flagsmith.com', + maxBuffer: 1, + flushInterval: 0); + p.trackEvent(event: 'auto'); + expect(p.buffer, isEmpty, reason: 'max buffer triggered a flush'); + expect(responses, 0); + + await p.flush(); + + expect(responses, 1, reason: 'flush must wait for the in-flight POST'); + }); + + test('When timer flush is in flight, then flush awaits its POST', () async { + var responses = 0; + dio.interceptors.add(InterceptorsWrapper(onResponse: (r, h) { + responses++; + h.next(r); + })); + adapter.onPost(eventsEndpoint, (server) { + server.reply(200, {}, + delay: const Duration(milliseconds: 80)); + }, data: Matchers.any); + + final p = processor(flushInterval: 20)..start(); + p.trackEvent(event: 'timed'); + await Future.delayed(const Duration(milliseconds: 40)); + expect(p.buffer, isEmpty, reason: 'timer flushed the buffer'); + expect(responses, 0); + + p.stop(); + await p.flush(); + + expect(responses, 1); + }); + + test('When stop is called, then timer cancelled and buffer flushed', + () async { + adapter.onPost(eventsEndpoint, (server) { + server.reply(200, {}); + }, data: Matchers.any); + + final p = processor(flushInterval: 20)..start(); + p.trackEvent(event: 'purchase'); + p.stop(); + await Future.delayed(const Duration(milliseconds: 10)); + expect(capture.requests.length, 1); + + p.trackEvent(event: 'after_stop'); + await Future.delayed(const Duration(milliseconds: 60)); + expect(capture.requests.length, 1, reason: 'timer must be cancelled'); + }); + }); +} diff --git a/test/shared.dart b/test/shared.dart index 4c8e9d2..31b4452 100644 --- a/test/shared.dart +++ b/test/shared.dart @@ -534,8 +534,99 @@ final flagsResponseData = r'''[ } ]'''; +final experimentFeatureName = 'experiment_enrolled'; +final experimentNotEnrolledFeatureName = 'experiment_not_enrolled'; +final experimentNoMetadataFeatureName = 'experiment_no_metadata'; +final experimentDisabledFeatureName = 'experiment_disabled'; +final experimentId = 42; +final experimentVariant = 'treatment-a'; + final identitiesResponseData = r'''{ "flags": [ + { + "id": 90001, + "feature": { + "id": 7001, + "name": "experiment_enrolled", + "created_date": "2026-09-01T08:38:29.203517Z", + "description": "Running experiment, identity enrolled", + "initial_value": null, + "default_enabled": false, + "type": "MULTIVARIATE" + }, + "feature_state_value": "buy-now", + "enabled": true, + "environment": 7822, + "identity": null, + "feature_segment": null, + "variant": "treatment-a", + "reason": "SPLIT; weight=50", + "metadata": { + "experiment": {"id": 42, "name": "New checkout CTA", "in_experiment": true}, + "unknown_key": {"ignored": true} + } + }, + { + "id": 90002, + "feature": { + "id": 7002, + "name": "experiment_not_enrolled", + "created_date": "2026-09-01T08:38:29.203517Z", + "description": "Running experiment, identity not enrolled", + "initial_value": null, + "default_enabled": false, + "type": "MULTIVARIATE" + }, + "feature_state_value": "control-value", + "enabled": true, + "environment": 7822, + "identity": null, + "feature_segment": null, + "variant": "control", + "reason": "DEFAULT", + "metadata": { + "experiment": {"id": 43, "name": "Not enrolled experiment", "in_experiment": false} + } + }, + { + "id": 90003, + "feature": { + "id": 7003, + "name": "experiment_no_metadata", + "created_date": "2026-09-01T08:38:29.203517Z", + "description": "Multivariate flag with no experiment", + "initial_value": null, + "default_enabled": false, + "type": "MULTIVARIATE" + }, + "feature_state_value": "control-value", + "enabled": true, + "environment": 7822, + "identity": null, + "feature_segment": null, + "variant": "control" + }, + { + "id": 90004, + "feature": { + "id": 7004, + "name": "experiment_disabled", + "created_date": "2026-09-01T08:38:29.203517Z", + "description": "Disabled flag inside a running experiment", + "initial_value": null, + "default_enabled": false, + "type": "MULTIVARIATE" + }, + "feature_state_value": null, + "enabled": false, + "environment": 7822, + "identity": null, + "feature_segment": null, + "variant": "treatment-a", + "metadata": { + "experiment": {"id": 44, "name": "Disabled experiment", "in_experiment": true} + } + }, { "id": 48540, "feature": {