diff --git a/CHANGES.txt b/CHANGES.txt index 3d681b60..aa1f7b44 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,3 +1,8 @@ +1.5.0 (June 29, 2022) +- Added a new config option to control the tasks that listen or poll for updates on feature flags and segments, via the new config sync.enabled . Running online Split will always pull the most recent updates upon initialization, this only affects updates fetching on a running instance. Useful when a consistent session experience is a must or to save resources when updates are not being used. +- Updated telemetry logic to track the anonymous config for user consent flag set to declined or unknown. +- Updated submitters logic, to avoid duplicating the post of impressions to Split cloud when the SDK is destroyed while its periodic post of impressions is running. + 1.4.1 (June 13, 2022) - Bugfixing - Updated submitters logic, to avoid dropping impressions and events that are being tracked while POST request is pending. diff --git a/package-lock.json b/package-lock.json index c959e174..2c6c6e5f 100644 --- a/package-lock.json +++ b/package-lock.json @@ -1,6 +1,6 @@ { "name": "@splitsoftware/splitio-commons", - "version": "1.4.1", + "version": "1.5.0", "lockfileVersion": 1, "requires": true, "dependencies": { @@ -4842,7 +4842,7 @@ "lodash.flatten": { "version": "4.4.0", "resolved": "https://registry.npmjs.org/lodash.flatten/-/lodash.flatten-4.4.0.tgz", - "integrity": "sha512-C5N2Z3DgnnKr0LOpv/hKCgKdb7ZZwafIrsesve6lmzvZIRZRGaZ/l6Q8+2W7NaT+ZwO3fFlSCzCzrDCFdJfZ4g==", + "integrity": "sha1-8xwiIlqWMtK7+OSt2+8kCqdlph8=", "dev": true }, "lodash.isarguments": { diff --git a/package.json b/package.json index d9951490..f7fb2df2 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "@splitsoftware/splitio-commons", - "version": "1.4.1", + "version": "1.5.0", "description": "Split Javascript SDK common components", "main": "cjs/index.js", "module": "esm/index.js", diff --git a/src/consent/__tests__/sdkUserConsent.spec.ts b/src/consent/__tests__/sdkUserConsent.spec.ts index 706b6ca9..e7981871 100644 --- a/src/consent/__tests__/sdkUserConsent.spec.ts +++ b/src/consent/__tests__/sdkUserConsent.spec.ts @@ -4,7 +4,7 @@ import { fullSettings } from '../../utils/settingsValidation/__tests__/settings. test('createUserConsentAPI', () => { const settings = { ...fullSettings, userConsent: 'UNKNOWN' }; - const syncManager = { submitter: syncTaskFactory() }; + const syncManager = { submitterManager: syncTaskFactory() }; const storage = { events: { clear: jest.fn() }, impressions: { clear: jest.fn() } @@ -20,15 +20,15 @@ test('createUserConsentAPI', () => { // setting user consent to 'GRANTED' expect(props.setStatus(true)).toBe(true); expect(props.setStatus(true)).toBe(true); // calling again has no affect - expect(syncManager.submitter.start).toBeCalledTimes(1); // submitter resumed - expect(syncManager.submitter.stop).toBeCalledTimes(0); + expect(syncManager.submitterManager.start).toBeCalledTimes(1); // submitter resumed + expect(syncManager.submitterManager.stop).toBeCalledTimes(0); expect(props.getStatus()).toBe(props.Status.GRANTED); // setting user consent to 'DECLINED' expect(props.setStatus(false)).toBe(true); expect(props.setStatus(false)).toBe(true); // calling again has no affect - expect(syncManager.submitter.start).toBeCalledTimes(1); - expect(syncManager.submitter.stop).toBeCalledTimes(1); // submitter paused + expect(syncManager.submitterManager.start).toBeCalledTimes(1); + expect(syncManager.submitterManager.stop).toBeCalledTimes(1); // submitter paused expect(props.getStatus()).toBe(props.Status.DECLINED); expect(storage.events.clear).toBeCalledTimes(1); // storage tracked data dropped expect(storage.impressions.clear).toBeCalledTimes(1); @@ -39,7 +39,7 @@ test('createUserConsentAPI', () => { expect(props.setStatus(undefined)).toBe(false); expect(props.setStatus({})).toBe(false); - expect(syncManager.submitter.start).toBeCalledTimes(1); - expect(syncManager.submitter.stop).toBeCalledTimes(1); + expect(syncManager.submitterManager.start).toBeCalledTimes(1); + expect(syncManager.submitterManager.stop).toBeCalledTimes(1); expect(props.getStatus()).toBe(props.Status.DECLINED); }); diff --git a/src/consent/sdkUserConsent.ts b/src/consent/sdkUserConsent.ts index ac8af3d8..e8f12156 100644 --- a/src/consent/sdkUserConsent.ts +++ b/src/consent/sdkUserConsent.ts @@ -34,9 +34,10 @@ export function createUserConsentAPI(params: ISdkFactoryContext) { settings.userConsent = newConsentStatus; if (consent) { // resumes submitters if transitioning to GRANTED - syncManager?.submitter?.start(); - } else { // pauses submitters and drops tracked data if transitioning to DECLINED - syncManager?.submitter?.stop(); + syncManager?.submitterManager?.start(); + } else { // pauses submitters (except telemetry), and drops tracked data if transitioning to DECLINED + syncManager?.submitterManager?.stop(true); + // @ts-ignore, clear method is present in storage for standalone and partial consumer mode if (events.clear) events.clear(); // @ts-ignore if (impressions.clear) impressions.clear(); diff --git a/src/listeners/__tests__/browser.spec.ts b/src/listeners/__tests__/browser.spec.ts index 566d6a68..b33bd709 100644 --- a/src/listeners/__tests__/browser.spec.ts +++ b/src/listeners/__tests__/browser.spec.ts @@ -258,11 +258,12 @@ test('Browser JS listener / standalone mode / user consent status', () => { settings.userConsent = 'DECLINED'; triggerUnloadEvent(); - // Unload event was triggered when user consent was unknown and declined. Thus sendBeacon and post services should not be called - expect(global.window.navigator.sendBeacon).toBeCalledTimes(0); + // Unload event was triggered when user consent was unknown and declined. Thus sendBeacon and post services should be called only for telemetry + expect(global.window.navigator.sendBeacon).toBeCalledTimes(2); expect(fakeSplitApi.postTestImpressionsBulk).not.toBeCalled(); expect(fakeSplitApi.postEventsBulk).not.toBeCalled(); expect(fakeSplitApi.postTestImpressionsCount).not.toBeCalled(); + (global.window.navigator.sendBeacon as jest.Mock).mockClear(); settings.userConsent = 'GRANTED'; triggerUnloadEvent(); diff --git a/src/listeners/browser.ts b/src/listeners/browser.ts index 77574eba..d50d4bed 100644 --- a/src/listeners/browser.ts +++ b/src/listeners/browser.ts @@ -67,7 +67,7 @@ export class BrowserSignalListener implements ISignalListener { flushData() { if (!this.syncManager) return; // In consumer mode there is not sync manager and data to flush - // Flush data if there is user consent + // Flush impressions & events data if there is user consent if (isConsentGranted(this.settings)) { const eventsUrl = this.settings.urls.events; const extraMetadata = { @@ -78,11 +78,13 @@ export class BrowserSignalListener implements ISignalListener { this._flushData(eventsUrl + '/testImpressions/beacon', this.storage.impressions, this.serviceApi.postTestImpressionsBulk, this.fromImpressionsCollector, extraMetadata); this._flushData(eventsUrl + '/events/beacon', this.storage.events, this.serviceApi.postEventsBulk); if (this.storage.impressionCounts) this._flushData(eventsUrl + '/testImpressions/count/beacon', this.storage.impressionCounts, this.serviceApi.postTestImpressionsCount, fromImpressionCountsCollector); - if (this.storage.telemetry) { - const telemetryUrl = this.settings.urls.telemetry; - const telemetryCacheAdapter = telemetryCacheStatsAdapter(this.storage.telemetry, this.storage.splits, this.storage.segments); - this._flushData(telemetryUrl + '/v1/metrics/usage/beacon', telemetryCacheAdapter, this.serviceApi.postMetricsUsage); - } + } + + // Flush telemetry data + if (this.storage.telemetry) { + const telemetryUrl = this.settings.urls.telemetry; + const telemetryCacheAdapter = telemetryCacheStatsAdapter(this.storage.telemetry, this.storage.splits, this.storage.segments); + this._flushData(telemetryUrl + '/v1/metrics/usage/beacon', telemetryCacheAdapter, this.serviceApi.postMetricsUsage); } // Close streaming connection diff --git a/src/sdkClient/clientAttributesDecoration.ts b/src/sdkClient/clientAttributesDecoration.ts index 5d55c7df..8d160108 100644 --- a/src/sdkClient/clientAttributesDecoration.ts +++ b/src/sdkClient/clientAttributesDecoration.ts @@ -5,7 +5,7 @@ import { ILogger } from '../logger/types'; import { objectAssign } from '../utils/lang/objectAssign'; /** - * Add in memory attributes storage methods and combine them with any attribute received from the getTreatment/s call + * Add in memory attributes storage methods and combine them with any attribute received from the getTreatment/s call */ export function clientAttributesDecoration(log: ILogger, client: TClient) { @@ -52,10 +52,10 @@ export function clientAttributesDecoration { @@ -101,8 +101,8 @@ export function clientAttributesDecoration { return { @@ -7,36 +9,173 @@ jest.mock('../submitters/submitterManager', () => { }; }); -import { syncManagerOnlineFactory } from '../syncManagerOnline'; +// Mocked storageManager +const storageManagerMock = { + splits: { + usesSegments: () => false + } +}; + +// @ts-expect-error +// Mocked readinessManager +let readinessManagerMock = { + isReady: jest.fn(() => true) // Fake the signal for the non ready SDK +} as IReadinessManager; + + +// Mocked pollingManager +const pollingManagerMock = { + syncAll: jest.fn(), + start: jest.fn(), + stop: jest.fn(), + isRunning: jest.fn(), + add: jest.fn(()=>{return {isrunning: () => true};}), + get: jest.fn() +}; + +const pushManagerMock = { + start: jest.fn(), + on: jest.fn(), + stop: jest.fn() +}; + +// Mocked pushManager +const pushManagerFactoryMock = jest.fn(() => pushManagerMock); test('syncManagerOnline should start or not the submitter depending on user consent status', () => { const settings = { ...fullSettings }; // @ts-ignore const syncManager = syncManagerOnlineFactory()({ settings }); - const submitter = syncManager.submitter!; + const submitterManager = syncManager.submitterManager!; syncManager.start(); - expect(submitter.start).toBeCalledTimes(1); // Submitter should be started if userConsent is undefined + expect(submitterManager.start).toBeCalledTimes(1); + expect(submitterManager.start).lastCalledWith(false); // SubmitterManager should start all submitters, if userConsent is undefined syncManager.stop(); - expect(submitter.stop).toBeCalledTimes(1); + expect(submitterManager.stop).toBeCalledTimes(1); settings.userConsent = 'UNKNOWN'; syncManager.start(); - expect(submitter.start).toBeCalledTimes(1); // Submitter should not be started if userConsent is unknown + expect(submitterManager.start).toBeCalledTimes(2); + expect(submitterManager.start).lastCalledWith(true); // SubmitterManager should start only telemetry submitter, if userConsent is unknown syncManager.stop(); - expect(submitter.stop).toBeCalledTimes(2); + expect(submitterManager.stop).toBeCalledTimes(2); + syncManager.flush(); + expect(submitterManager.execute).toBeCalledTimes(1); + expect(submitterManager.execute).lastCalledWith(true); // SubmitterManager should flush only telemetry, if userConsent is unknown settings.userConsent = 'GRANTED'; syncManager.start(); - expect(submitter.start).toBeCalledTimes(2); // Submitter should be started if userConsent is granted + expect(submitterManager.start).toBeCalledTimes(3); + expect(submitterManager.start).lastCalledWith(false); // SubmitterManager should start all submitters, if userConsent is granted syncManager.stop(); - expect(submitter.stop).toBeCalledTimes(3); + expect(submitterManager.stop).toBeCalledTimes(3); + syncManager.flush(); + expect(submitterManager.execute).toBeCalledTimes(2); + expect(submitterManager.execute).lastCalledWith(false); // SubmitterManager should flush all submitters, if userConsent is granted settings.userConsent = 'DECLINED'; syncManager.start(); - expect(submitter.start).toBeCalledTimes(2); // Submitter should not be started if userConsent is declined + expect(submitterManager.start).toBeCalledTimes(4); + expect(submitterManager.start).lastCalledWith(true); // SubmitterManager should start only telemetry submitter, if userConsent is declined + syncManager.stop(); + expect(submitterManager.stop).toBeCalledTimes(4); + syncManager.flush(); + expect(submitterManager.execute).toBeCalledTimes(3); + expect(submitterManager.execute).lastCalledWith(true); // SubmitterManager should flush only telemetry, if userConsent is unknown + +}); + +test('syncManagerOnline should syncAll a single time when sync is disabled', () => { + const settings = { ...fullSettings }; + + // disable sync + settings.sync.enabled = false; + + // @ts-ignore + // Test pushManager for main client + const syncManager = syncManagerOnlineFactory(() => pollingManagerMock, pushManagerFactoryMock)({ settings }); + + expect(pushManagerFactoryMock).not.toBeCalled(); + + // Test pollingManager for Main client + syncManager.start(); + + expect(pollingManagerMock.start).not.toBeCalled(); + expect(pollingManagerMock.syncAll).toBeCalledTimes(1); + syncManager.stop(); - expect(submitter.stop).toBeCalledTimes(4); + syncManager.start(); + + expect(pollingManagerMock.start).not.toBeCalled(); + expect(pollingManagerMock.syncAll).toBeCalledTimes(1); + + syncManager.stop(); + syncManager.start(); + + expect(pollingManagerMock.start).not.toBeCalled(); + expect(pollingManagerMock.syncAll).toBeCalledTimes(1); + + syncManager.stop(); + + // @ts-ignore + // Test pollingManager for shared client + const pollingSyncManagerShared = syncManager.shared('sharedKey', readinessManagerMock, storageManagerMock); + + if (!pollingSyncManagerShared) throw new Error('pollingSyncManagerShared should exist'); + + pollingSyncManagerShared.start(); + + expect(pollingManagerMock.start).not.toBeCalled(); + + pollingSyncManagerShared.stop(); + pollingSyncManagerShared.start(); + + expect(pollingManagerMock.start).not.toBeCalled(); + + pollingSyncManagerShared.stop(); + + syncManager.start(); + + expect(pollingManagerMock.start).not.toBeCalled(); + + syncManager.stop(); + syncManager.start(); + + expect(pollingManagerMock.start).not.toBeCalled(); + + syncManager.stop(); + + // @ts-ignore + // Test pollingManager for shared client + const pushingSyncManagerShared = syncManager.shared('pushingSharedKey', readinessManagerMock, storageManagerMock); + + if (!pushingSyncManagerShared) throw new Error('pushingSyncManagerShared should exist'); + + pushingSyncManagerShared.start(); + + expect(pollingManagerMock.start).not.toBeCalled(); + + pushingSyncManagerShared.stop(); + pushingSyncManagerShared.start(); + + expect(pollingManagerMock.start).not.toBeCalled(); + + pushingSyncManagerShared.stop(); + + settings.sync.enabled = true; + // @ts-ignore + // pushManager instantiation control test + const testSyncManager = syncManagerOnlineFactory(() => pollingManagerMock, pushManagerFactoryMock)({ settings }); + + expect(pushManagerFactoryMock).toBeCalled(); + + // Test pollingManager for Main client + testSyncManager.start(); + + expect(pushManagerMock.start).toBeCalled(); + + testSyncManager.stop(); }); diff --git a/src/sync/__tests__/syncTask.spec.ts b/src/sync/__tests__/syncTask.spec.ts index 49ecf67f..22586e5a 100644 --- a/src/sync/__tests__/syncTask.spec.ts +++ b/src/sync/__tests__/syncTask.spec.ts @@ -5,7 +5,7 @@ const period = 30; const taskResult = 'taskResult'; const asyncTask = jest.fn(() => Promise.resolve(taskResult)); -test('syncTaskFactory', (done) => { +test('syncTaskFactory / start & stop methods for periodic execution', async () => { const syncTask = syncTaskFactory(loggerMock, asyncTask, period); @@ -40,45 +40,75 @@ test('syncTaskFactory', (done) => { expect(syncTask.isExecuting()).toBe(true); // Executing - setTimeout(() => { - expect(asyncTask).toHaveBeenLastCalledWith(...startArgs); // Periodic call should be done with the initial `start` arguments - expect(asyncTask).toBeCalledTimes(4); // The task was executed 4 times: twice due to periodic execution and twice due to execute call - - setTimeout(() => { - expect(asyncTask).toHaveBeenLastCalledWith(...startArgs); // Periodic call should be done with the initial `start` arguments - expect(asyncTask).toBeCalledTimes(5); // The task was executed 5 times: 3 due to periodic execution and twice due to execute call - - // Calling `stop` stops the periodic execution of the given task - expect(syncTask.isRunning()).toBe(true); // Running periodically - syncTask.stop(); - expect(syncTask.isRunning()).toBe(false); // Stop running periodically - - setTimeout(() => { - expect(asyncTask).toBeCalledTimes(5); // Stopped task should not be called again - - // Stopping and starting - syncTask.stop(); - syncTask.start(); // Inmediatelly call task - syncTask.stop(); - syncTask.start(); // Inmediatelly call task - syncTask.stop(); - expect(asyncTask).toBeCalledTimes(7); - - // Resume periodic execution - syncTask.start(); // Inmediatelly call task - syncTask.start(); // No effect - expect(asyncTask).toBeCalledTimes(8); - - setTimeout(() => { - expect(asyncTask).toBeCalledTimes(9); // Stopped task should not be called again - expect(syncTask.isRunning()).toBe(true); // Running periodically - syncTask.stop(); // Finally stop to finish the test - expect(syncTask.isRunning()).toBe(false); // Stop running periodically - - done(); - }, period + 10); - }, period + 10); - }, period + 10); - }, period + 10); + await new Promise(res => setTimeout(res, period + 10)); + expect(asyncTask).toHaveBeenLastCalledWith(...startArgs); // Periodic call should be done with the initial `start` arguments + expect(asyncTask).toBeCalledTimes(4); // The task was executed 4 times: twice due to periodic execution and twice due to execute call + + await new Promise(res => setTimeout(res, period + 10)); + expect(asyncTask).toHaveBeenLastCalledWith(...startArgs); // Periodic call should be done with the initial `start` arguments + expect(asyncTask).toBeCalledTimes(5); // The task was executed 5 times: 3 due to periodic execution and twice due to execute call + + // Calling `stop` stops the periodic execution of the given task + expect(syncTask.isRunning()).toBe(true); // Running periodically + syncTask.stop(); + expect(syncTask.isRunning()).toBe(false); // Stop running periodically + + await new Promise(res => setTimeout(res, period + 10)); + expect(asyncTask).toBeCalledTimes(5); // Stopped task should not be called again + + // Stopping and starting + syncTask.stop(); + syncTask.start(); // Inmediatelly call task + syncTask.stop(); + syncTask.start(); // Doesn't call task since previous one has not been resolved + syncTask.stop(); + expect(asyncTask).toBeCalledTimes(6); + + // Resume periodic execution + syncTask.start(); // Doesn't call task since previous one has not been resolved + syncTask.start(); // No effect + expect(asyncTask).toBeCalledTimes(6); + + await new Promise(res => setTimeout(res, period + 10)); + expect(asyncTask).toBeCalledTimes(9); // Stopped task should not be called again + expect(syncTask.isRunning()).toBe(true); // Running periodically + syncTask.stop(); // Finally stop to finish the test + expect(syncTask.isRunning()).toBe(false); // Stop running periodically + +}); + +test('syncTaskFactory / execute method', (done) => { + let executeCount = 0; + let resolveOrder = 0; + const asyncTask = jest.fn((toReturn) => { + executeCount++; + return new Promise(res => setTimeout(() => res(toReturn))); + }); + + const syncTask = syncTaskFactory(loggerMock, asyncTask, period); + + syncTask.execute(1).then(result=> { + // console.log('1'); + resolveOrder++; + expect(resolveOrder).toBe(1); // @TODO should be 1 ? + expect(executeCount).toBe(1); + expect(result).toBe(1); // @TODO should be 1 + }); + syncTask.execute(2).then(result=> { + // console.log('2'); + resolveOrder++; + expect(resolveOrder).toBe(2); + expect(executeCount).toBe(3); // @TODO borrar + expect(result).toBe(2); // @TODO should be 2 + }); + syncTask.execute(3).then(result=> { + // console.log('3'); + resolveOrder++; + expect(resolveOrder).toBe(3); // @TODO should be 3 ? + expect(executeCount).toBe(3); + expect(result).toBe(3); + + done(); + }); }); diff --git a/src/sync/submitters/__tests__/eventsSubmitter.spec.ts b/src/sync/submitters/__tests__/eventsSubmitter.spec.ts index d8336b4f..d79d0c5e 100644 --- a/src/sync/submitters/__tests__/eventsSubmitter.spec.ts +++ b/src/sync/submitters/__tests__/eventsSubmitter.spec.ts @@ -48,7 +48,7 @@ describe('Events submitter', () => { expect(eventsSubmitter.isRunning()).toEqual(false); }); - test('without eventsFirstPushWindow', async () => { + test('without eventsFirstPushWindow', (done) => { const eventsFirstPushWindow = 0; params.settings.startup.eventsFirstPushWindow = eventsFirstPushWindow; // @ts-ignore const eventsSubmitter = eventsSubmitterFactory(params); @@ -58,14 +58,19 @@ describe('Events submitter', () => { expect(eventsSubmitter.isExecuting()).toEqual(true); // and executes immediatelly if there isn't a push window expect(eventsCacheMock.isEmpty).toBeCalledTimes(1); - // If queue is full, submitter should be executed + // If queue is full, submitter is executed again after current execution is resolved __onFullQueueCb(); - expect(eventsSubmitter.isExecuting()).toEqual(true); - expect(eventsCacheMock.isEmpty).toBeCalledTimes(2); + expect(eventsCacheMock.isEmpty).toBeCalledTimes(1); - expect(eventsSubmitter.isRunning()).toEqual(true); - eventsSubmitter.stop(); - expect(eventsSubmitter.isRunning()).toEqual(false); + setTimeout(()=> { + expect(eventsSubmitter.isExecuting()).toEqual(false); + expect(eventsCacheMock.isEmpty).toBeCalledTimes(2); // 2 executions: 1st due to start and 2nd due to full queue + + expect(eventsSubmitter.isRunning()).toEqual(true); + eventsSubmitter.stop(); + expect(eventsSubmitter.isRunning()).toEqual(false); + done(); + }); }); test('doesn\'t drop items from cache when POST is resolved', (done) => { diff --git a/src/sync/submitters/__tests__/impressionsSubmitter.spec.ts b/src/sync/submitters/__tests__/impressionsSubmitter.spec.ts index 4f59fb19..cffd1b59 100644 --- a/src/sync/submitters/__tests__/impressionsSubmitter.spec.ts +++ b/src/sync/submitters/__tests__/impressionsSubmitter.spec.ts @@ -74,4 +74,26 @@ describe('Impressions submitter', () => { }, params.settings.scheduler.impressionsPushRate + 10); }); + test('if it is executed while POST is pending, execution is queued until POST is resolved and not same items are submitted', (done) => { + // Make the POST request fail + params.splitApi.postTestImpressionsBulk.mockImplementation(() => Promise.resolve()); + + impressionsCacheInMemory.track([imp1]); + impressionsSubmitter.start(); + + // Tracking impression and executing submitter while POST is pending + impressionsCacheInMemory.track([{ ...imp1, keyName: 'k2' }]); + impressionsSubmitter.execute().then(() => { + expect(params.splitApi.postTestImpressionsBulk.mock.calls).toEqual([ + // impression for k1 + ['[{"f":"someFeature","i":[{"k":"k1","t":"someTreatment","m":0,"c":123}]}]'], + // impression for k2 + ['[{"f":"someFeature","i":[{"k":"k2","t":"someTreatment","m":0,"c":123}]}]']]); + impressionsSubmitter.stop(); + + done(); + }); + + }); + }); diff --git a/src/sync/submitters/submitterManager.ts b/src/sync/submitters/submitterManager.ts index 523e5ab5..298f61a4 100644 --- a/src/sync/submitters/submitterManager.ts +++ b/src/sync/submitters/submitterManager.ts @@ -1,11 +1,11 @@ -import { syncTaskComposite } from '../syncTaskComposite'; import { eventsSubmitterFactory } from './eventsSubmitter'; import { impressionsSubmitterFactory } from './impressionsSubmitter'; import { impressionCountsSubmitterFactory } from './impressionCountsSubmitter'; import { telemetrySubmitterFactory } from './telemetrySubmitter'; import { ISdkFactoryContextSync } from '../../sdkFactory/types'; +import { ISubmitterManager } from './types'; -export function submitterManagerFactory(params: ISdkFactoryContextSync) { +export function submitterManagerFactory(params: ISdkFactoryContextSync): ISubmitterManager { const submitters = [ impressionsSubmitterFactory(params), @@ -15,7 +15,33 @@ export function submitterManagerFactory(params: ISdkFactoryContextSync) { const impressionCountsSubmitter = impressionCountsSubmitterFactory(params); if (impressionCountsSubmitter) submitters.push(impressionCountsSubmitter); const telemetrySubmitter = telemetrySubmitterFactory(params); - if (telemetrySubmitter) submitters.push(telemetrySubmitter); - return syncTaskComposite(submitters); + return { + // `onlyTelemetry` true if SDK is created with userConsent not GRANTED + start(onlyTelemetry?: boolean) { + if (!onlyTelemetry) submitters.forEach(submitter => submitter.start()); + if (telemetrySubmitter) telemetrySubmitter.start(); + }, + + // `allExceptTelemetry` true if userConsent is changed to DECLINED + stop(allExceptTelemetry?: boolean) { + submitters.forEach(submitter => submitter.stop()); + if (!allExceptTelemetry && telemetrySubmitter) telemetrySubmitter.stop(); + }, + + isRunning() { + return submitters.some(submitter => submitter.isRunning()); + }, + + // Flush data. Called with `onlyTelemetry` true if SDK is destroyed with userConsent not GRANTED + execute(onlyTelemetry?: boolean) { + const promises = onlyTelemetry ? [] : submitters.map(submitter => submitter.execute()); + if (telemetrySubmitter) promises.push(telemetrySubmitter.execute()); + return Promise.all(promises); + }, + + isExecuting() { + return submitters.some(submitter => submitter.isExecuting()); + } + }; } diff --git a/src/sync/submitters/types.ts b/src/sync/submitters/types.ts index f4eb8c7b..fe9fac26 100644 --- a/src/sync/submitters/types.ts +++ b/src/sync/submitters/types.ts @@ -1,5 +1,6 @@ import { IMetadata } from '../../dtos/types'; import { SplitIO } from '../../types'; +import { ISyncTask } from '../types'; export type ImpressionsPayload = { /** Split name */ @@ -191,3 +192,9 @@ export type TelemetryConfigStatsPayload = TelemetryConfigStats & { i?: Array, // integrations uC: number, // userConsent } + +export interface ISubmitterManager extends ISyncTask { + start(onlyTelemetry?: boolean): void, + stop(allExceptTelemetry?: boolean): void, + execute(onlyTelemetry?: boolean): Promise +} diff --git a/src/sync/syncManagerOnline.ts b/src/sync/syncManagerOnline.ts index 61f0603d..b6407630 100644 --- a/src/sync/syncManagerOnline.ts +++ b/src/sync/syncManagerOnline.ts @@ -28,19 +28,19 @@ export function syncManagerOnlineFactory( */ return function (params: ISdkFactoryContextSync): ISyncManagerCS { - const { settings, settings: { log, streamingEnabled }, telemetryTracker } = params; + const { settings, settings: { log, streamingEnabled, sync: { enabled: syncEnabled } }, telemetryTracker } = params; /** Polling Manager */ const pollingManager = pollingManagerFactory && pollingManagerFactory(params); /** Push Manager */ - const pushManager = streamingEnabled && pollingManager && pushManagerFactory ? + const pushManager = syncEnabled && streamingEnabled && pollingManager && pushManagerFactory ? pushManagerFactory(params, pollingManager) : undefined; /** Submitter Manager */ // It is not inyected as push and polling managers, because at the moment it is required - const submitter = submitterManagerFactory(params); + const submitterManager = submitterManagerFactory(params); /** Sync Manager logic */ @@ -79,7 +79,7 @@ export function syncManagerOnlineFactory( // E.g.: user consent, app state changes (Page hide, Foreground/Background, Online/Offline). pollingManager, pushManager, - submitter, + submitterManager, /** * Method used to start the syncManager for the first time, or resume it after being stopped. @@ -89,20 +89,29 @@ export function syncManagerOnlineFactory( // start syncing splits and segments if (pollingManager) { - if (pushManager) { - // Doesn't call `syncAll` when the syncManager is resuming + + // If synchronization is disabled pushManager and pollingManager should not start + if (syncEnabled) { + if (pushManager) { + // Doesn't call `syncAll` when the syncManager is resuming + if (startFirstTime) { + pollingManager.syncAll(); + startFirstTime = false; + } + pushManager.start(); + } else { + pollingManager.start(); + } + } else { if (startFirstTime) { pollingManager.syncAll(); startFirstTime = false; } - pushManager.start(); - } else { - pollingManager.start(); } } // start periodic data recording (events, impressions, telemetry). - if (isConsentGranted(settings)) submitter.start(); + submitterManager.start(!isConsentGranted(settings)); }, /** @@ -116,7 +125,7 @@ export function syncManagerOnlineFactory( if (pollingManager && pollingManager.isRunning()) pollingManager.stop(); // stop periodic data recording (events, impressions, telemetry). - submitter.stop(); + submitterManager.stop(); }, isRunning() { @@ -124,8 +133,7 @@ export function syncManagerOnlineFactory( }, flush() { - if (isConsentGranted(settings)) return submitter.execute(); - else return Promise.resolve(); + return submitterManager.execute(!isConsentGranted(settings)); }, // [Only used for client-side] @@ -138,18 +146,22 @@ export function syncManagerOnlineFactory( return { isRunning: mySegmentsSyncTask.isRunning, start() { - if (pushManager) { - if (pollingManager!.isRunning()) { - // if doing polling, we must start the periodic fetch of data - if (storage.splits.usesSegments()) mySegmentsSyncTask.start(); + if (syncEnabled) { + if (pushManager) { + if (pollingManager!.isRunning()) { + // if doing polling, we must start the periodic fetch of data + if (storage.splits.usesSegments()) mySegmentsSyncTask.start(); + } else { + // if not polling, we must execute the sync task for the initial fetch + // of segments since `syncAll` was already executed when starting the main client + mySegmentsSyncTask.execute(); + } + pushManager.add(matchingKey, mySegmentsSyncTask); } else { - // if not polling, we must execute the sync task for the initial fetch - // of segments since `syncAll` was already executed when starting the main client - mySegmentsSyncTask.execute(); + if (storage.splits.usesSegments()) mySegmentsSyncTask.start(); } - pushManager.add(matchingKey, mySegmentsSyncTask); } else { - if (storage.splits.usesSegments()) mySegmentsSyncTask.start(); + if (!readinessManager.isReady()) mySegmentsSyncTask.execute(); } }, stop() { diff --git a/src/sync/syncTask.ts b/src/sync/syncTask.ts index eb3a2c30..0c6abf73 100644 --- a/src/sync/syncTask.ts +++ b/src/sync/syncTask.ts @@ -4,8 +4,8 @@ import { ISyncTask } from './types'; /** * Creates a syncTask that handles the periodic execution of a given task ("start" and "stop" methods). - * The task can be executed once calling the "execute" method. - * NOTE: Multiple calls to "execute" are not queued. Use "isExecuting" method to handle synchronization. + * The task can be also executed by calling the "execute" method. Multiple execute calls are chained to run secuentially and avoid race conditions. + * For example, submitters executed on SDK destroy or full queue, while periodic execution is pending. * * @param log Logger instance. * @param task Task to execute that returns a promise that NEVER REJECTS. Otherwise, periodic execution can result in Unhandled Promise Rejections. @@ -15,8 +15,8 @@ import { ISyncTask } from './types'; */ export function syncTaskFactory(log: ILogger, task: (...args: Input) => Promise, period: number, taskName = 'task'): ISyncTask { - // Flag that indicates if the task is being executed - let executing = false; + // Task promise while it is pending. Undefined once the promise is resolved + let pendingTask: Promise | undefined; // flag that indicates if the task periodic execution has been started/stopped. let running = false; // Auxiliar counter used to avoid race condition when calling `start` & `stop` intermittently @@ -26,14 +26,21 @@ export function syncTaskFactory(log: ILogger, // Id of the periodic call timeout let timeoutID: any; - function execute(...args: Input) { - executing = true; + function execute(...args: Input): Promise { + // If task is executing, chain the new execution + if (pendingTask) { + return pendingTask.then(() => { + return execute(...args); + }); + } + + // Execute task log.debug(SYNC_TASK_EXECUTE, [taskName]); - return task(...args).then(result => { - executing = false; + pendingTask = task(...args).then(result => { + pendingTask = undefined; return result; }); - // No need to handle promise rejection because it is a pre-condition that provided task never rejects. + return pendingTask; } function periodicExecute(currentRunningId: number) { @@ -46,11 +53,10 @@ export function syncTaskFactory(log: ILogger, } return { - // @TODO check if we need to queued `execute` calls, to avoid possible race conditions on submitters and updaters with streaming. execute, isExecuting() { - return executing; + return pendingTask !== undefined; }, start(...args: Input) { diff --git a/src/sync/syncTaskComposite.ts b/src/sync/syncTaskComposite.ts deleted file mode 100644 index b16a1704..00000000 --- a/src/sync/syncTaskComposite.ts +++ /dev/null @@ -1,26 +0,0 @@ -import { ISyncTask } from './types'; - -/** - * Composite Sync Task: group of sync tasks that are treated as a single one. - */ -export function syncTaskComposite(syncTasks: ISyncTask[]): ISyncTask { - - return { - start() { - syncTasks.forEach(syncTask => syncTask.start()); - }, - stop() { - syncTasks.forEach(syncTask => syncTask.stop()); - }, - isRunning() { - return syncTasks.some(syncTask => syncTask.isRunning()); - }, - execute() { - return Promise.all(syncTasks.map(syncTask => syncTask.execute())); - }, - isExecuting() { - return syncTasks.some(syncTask => syncTask.isExecuting()); - } - }; - -} diff --git a/src/sync/types.ts b/src/sync/types.ts index 22f51eb9..81727ca9 100644 --- a/src/sync/types.ts +++ b/src/sync/types.ts @@ -2,6 +2,7 @@ import { IReadinessManager } from '../readiness/types'; import { IStorageSync } from '../storages/types'; import { IPollingManager } from './polling/types'; import { IPushManager } from './streaming/types'; +import { ISubmitterManager } from './submitters/types'; export interface ITask { /** @@ -39,7 +40,7 @@ export interface ISyncManager extends ITask { flush(): Promise, pushManager?: IPushManager, pollingManager?: IPollingManager, - submitter?: ISyncTask + submitterManager?: ISubmitterManager } export interface ISyncManagerCS extends ISyncManager { diff --git a/src/types.ts b/src/types.ts index dd2cb085..a45c27a7 100644 --- a/src/types.ts +++ b/src/types.ts @@ -117,7 +117,8 @@ export interface ISettings { splitFilters: SplitIO.SplitFilter[], impressionsMode: SplitIO.ImpressionsMode, __splitFiltersValidation: ISplitFiltersValidation, - localhostMode?: SplitIO.LocalhostFactory + localhostMode?: SplitIO.LocalhostFactory, + enabled: boolean }, readonly runtime: { ip: string | false @@ -214,6 +215,11 @@ interface ISharedSettings { * @default 'OPTIMIZED' */ impressionsMode?: SplitIO.ImpressionsMode, + /** + * Enables synchronization. + * @property {boolean} enabled + */ + enabled: boolean } } /** diff --git a/src/utils/settingsValidation/__tests__/index.spec.ts b/src/utils/settingsValidation/__tests__/index.spec.ts index e82d1d99..2b3f4918 100644 --- a/src/utils/settingsValidation/__tests__/index.spec.ts +++ b/src/utils/settingsValidation/__tests__/index.spec.ts @@ -40,6 +40,7 @@ describe('settingsValidation', () => { telemetry: 'https://telemetry.split.io/api', }); expect(settings.sync.impressionsMode).toBe(OPTIMIZED); + expect(settings.sync.enabled).toBe(true); }); test('override with default impressionMode if provided one is invalid', () => { @@ -163,6 +164,27 @@ describe('settingsValidation', () => { expect(settingsWithStreamingEnabled.streamingEnabled).toBe(true); // If streamingEnabled is not provided, it will be true. }); + test('sync.enabled should be overwritable and true by default', () => { + const settingsWithSyncEnabled = settingsValidation({ + core: { + authorizationKey: 'dummy token', + } + }, minimalSettingsParams); + + const settingsWithSyncDisabled = settingsValidation({ + core: { + authorizationKey: 'dummy token' + }, + sync: { + enabled: false + } + }, minimalSettingsParams); + + expect(settingsWithSyncDisabled.sync.enabled).toBe(false); // If sync.enabled is not provided, it will be true. + expect(settingsWithSyncEnabled.sync.enabled).toBe(true); // When creating a setting instance, it will have the provided value for sync.enabled + + }); + const storageMock = () => { }; const integrationsMock = [() => { }]; diff --git a/src/utils/settingsValidation/__tests__/settings.mocks.ts b/src/utils/settingsValidation/__tests__/settings.mocks.ts index 908bfdc4..8d213b03 100644 --- a/src/utils/settingsValidation/__tests__/settings.mocks.ts +++ b/src/utils/settingsValidation/__tests__/settings.mocks.ts @@ -78,7 +78,8 @@ export const fullSettings: ISettings = { validFilters: [], queryString: null, groupedFilters: { byName: [], byPrefix: [] } - } + }, + enabled: true }, version: 'jest', runtime: { diff --git a/src/utils/settingsValidation/index.ts b/src/utils/settingsValidation/index.ts index 2269721f..84d97954 100644 --- a/src/utils/settingsValidation/index.ts +++ b/src/utils/settingsValidation/index.ts @@ -83,7 +83,8 @@ export const base = { splitFilters: undefined, // impressions collection mode impressionsMode: OPTIMIZED, - localhostMode: undefined + localhostMode: undefined, + enabled: true }, // Logger @@ -191,6 +192,11 @@ export function settingsValidation(config: unknown, validationParams: ISettingsV scheduler.pushRetryBackoffBase = fromSecondsToMillis(scheduler.pushRetryBackoffBase); } + // validate sync enabled + if (withDefaults.sync.enabled !== false) { // @ts-ignore, modify readonly prop + withDefaults.sync.enabled = true; + } + // validate the `splitFilters` settings and parse splits query const splitFiltersValidation = validateSplitFilters(log, withDefaults.sync.splitFilters, withDefaults.mode); withDefaults.sync.splitFilters = splitFiltersValidation.validFilters;