Skip to content
Merged
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
26 changes: 14 additions & 12 deletions src/services/splitApi.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,8 @@ import { splitHttpClientFactory } from './splitHttpClient';
import { ISplitApi } from './types';
import { ISettingsInternal } from '../utils/settingsValidation/types';

const noCacheHeaderOptions = { headers: { 'Cache-Control': 'no-cache' } };

function userKeyToQueryParam(userKey: string) {
return 'users=' + encodeURIComponent(userKey); // no need to check availability of `encodeURIComponent`, since it is a global highly supported.
}
Expand Down Expand Up @@ -34,53 +36,53 @@ export function splitApiFactory(settings: ISettings, platform: IPlatform): ISpli
return splitHttpClient(url);
},

fetchSplitChanges(since: number) {
fetchSplitChanges(since: number, noCache?: boolean) {
const url = `${urls.sdk}/splitChanges?since=${since}${filterQueryString || ''}`;
return splitHttpClient(url);
return splitHttpClient(url, noCache ? noCacheHeaderOptions : undefined);
},

fetchSegmentChanges(since: number, segmentName: string) {
fetchSegmentChanges(since: number, segmentName: string, noCache?: boolean) {
const url = `${urls.sdk}/segmentChanges/${segmentName}?since=${since}`;
return splitHttpClient(url);
return splitHttpClient(url, noCache ? noCacheHeaderOptions : undefined);
},

fetchMySegments(userMatchingKey: string) {
fetchMySegments(userMatchingKey: string, noCache?: boolean) {
/**
* URI encoding of user keys in order to:
* - avoid 400 responses (due to URI malformed). E.g.: '/api/mySegments/%'
* - avoid 404 responses. E.g.: '/api/mySegments/foo/bar'
* - match user keys with special characters. E.g.: 'foo%bar', 'foo/bar'
*/
const url = `${urls.sdk}/mySegments/${encodeURIComponent(userMatchingKey)}`;
return splitHttpClient(url);
return splitHttpClient(url, noCache ? noCacheHeaderOptions : undefined);
},

postEventsBulk(body: string) {
const url = `${urls.events}/events/bulk`;
return splitHttpClient(url, 'POST', body);
return splitHttpClient(url, { method: 'POST', body });
},

postTestImpressionsBulk(body: string) {
const url = `${urls.events}/testImpressions/bulk`;
return splitHttpClient(url, 'POST', body, false, {
return splitHttpClient(url, {
// Adding extra headers to send impressions in OPTIMIZED or DEBUG modes.
SplitSDKImpressionsMode
method: 'POST', body, headers: { SplitSDKImpressionsMode }
});
},

postTestImpressionsCount(body: string) {
const url = `${urls.events}/testImpressions/count`;
return splitHttpClient(url, 'POST', body);
return splitHttpClient(url, { method: 'POST', body });
},

postMetricsCounters(body: string) {
const url = `${urls.events}/metrics/counters`;
return splitHttpClient(url, 'POST', body, true);
return splitHttpClient(url, { method: 'POST', body }, true);
},

postMetricsTimes(body: string) {
const url = `${urls.events}/metrics/times`;
return splitHttpClient(url, 'POST', body, true);
return splitHttpClient(url, { method: 'POST', body }, true);
}
};
}
12 changes: 8 additions & 4 deletions src/services/splitHttpClient.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
import { IFetch, ISplitHttpClient } from './types';
import { IFetch, IRequestOptions, ISplitHttpClient } from './types';
import { SplitError, SplitNetworkError } from '../utils/lang/errors';
import objectAssign from 'object-assign';
import { logFactory } from '../logger/sdkLogger';
Expand Down Expand Up @@ -33,9 +33,13 @@ export function splitHttpClientFactory(apikey: string, metadata: IMetadata, getF
if (metadata.ip) headers['SplitSDKMachineIP'] = metadata.ip;
if (metadata.hostname) headers['SplitSDKMachineName'] = metadata.hostname;

return function httpClient(url: string, method: string = 'GET', body?: string, logErrorsAsInfo: boolean = false, extraHeaders?: Record<string, string>): Promise<Response> {
const rHeaders = extraHeaders ? objectAssign({}, headers, extraHeaders) : headers;
const request = objectAssign({ headers: rHeaders, method, body }, options);
return function httpClient(url: string, reqOpts: IRequestOptions = {}, logErrorsAsInfo: boolean = false): Promise<Response> {

const request = objectAssign({
headers: reqOpts.headers ? objectAssign({}, headers, reqOpts.headers) : headers,
method: reqOpts.method || 'GET',
body: reqOpts.body
}, options);

// using `fetch(url, options)` signature to work with unfetch, a lightweight ponyfill of fetch API.
return fetch ? fetch(url, request)
Expand Down
27 changes: 13 additions & 14 deletions src/services/types.ts
Original file line number Diff line number Diff line change
@@ -1,22 +1,21 @@
export type IFetch = (
url: string,
options?: {
method?: string,
headers?: Record<string, string>,
credentials?: 'include' | 'omit',
body?: string
}
) => Promise<Response>

export type ISplitHttpClient = (url: string, method?: string, body?: string, logErrorsAsInfo?: boolean, extraHeaders?: Record<string, string>) => Promise<Response>
export type IRequestOptions = {
method?: string,
headers?: Record<string, string>,
body?: string
};

// Reduced version of Fetch API
export type IFetch = (url: string, options?: IRequestOptions) => Promise<Response>

export type ISplitHttpClient = (url: string, options?: IRequestOptions, logErrorsAsInfo?: boolean) => Promise<Response>

export type IFetchAuth = (userKeys?: string[]) => Promise<Response>

export type IFetchSplitChanges = (since: number) => Promise<Response>
export type IFetchSplitChanges = (since: number, noCache?: boolean) => Promise<Response>

export type IFetchSegmentChanges = (since: number, segmentName: string) => Promise<Response>
export type IFetchSegmentChanges = (since: number, segmentName: string, noCache?: boolean) => Promise<Response>

export type IFetchMySegments = (userMatchingKey: string) => Promise<Response>
export type IFetchMySegments = (userMatchingKey: string, noCache?: boolean) => Promise<Response>

export type IPostEventsBulk = (body: string) => Promise<Response>

Expand Down
3 changes: 2 additions & 1 deletion src/sync/polling/fetchers/mySegmentsFetcher.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,11 +10,12 @@ import { IMySegmentsFetcher } from './types';
export default function mySegmentsFetcherFactory(fetchMySegments: IFetchMySegments, userMatchingKey: string): IMySegmentsFetcher {

return function mySegmentsFetcher(
noCache?: boolean,
// Optional decorator for `fetchMySegments` promise, such as timeout or time tracker
decorator?: (promise: Promise<Response>) => Promise<Response>
): Promise<string[]> {

let mySegmentsPromise = fetchMySegments(userMatchingKey);
let mySegmentsPromise = fetchMySegments(userMatchingKey, noCache);
if (decorator) mySegmentsPromise = decorator(mySegmentsPromise);

// Extract segment names
Expand Down
9 changes: 5 additions & 4 deletions src/sync/polling/fetchers/segmentChangesFetcher.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,16 +2,16 @@ import { IFetchSegmentChanges } from '../../../services/types';
import { ISegmentChangesResponse } from '../../../dtos/types';
import { ISegmentChangesFetcher } from './types';

function greedyFetch(fetchSegmentChanges: IFetchSegmentChanges, since: number, segmentName: string): Promise<ISegmentChangesResponse[]> {
return fetchSegmentChanges(since, segmentName)
function greedyFetch(fetchSegmentChanges: IFetchSegmentChanges, since: number, segmentName: string, noCache?: boolean): Promise<ISegmentChangesResponse[]> {
return fetchSegmentChanges(since, segmentName, noCache)
// no need to handle json parsing errors as SplitError, since errors are handled differently for segments
.then(resp => resp.json())
.then((json: ISegmentChangesResponse) => {
let { since, till } = json;
if (since === till) {
return [json];
} else {
return Promise.all([json, greedyFetch(fetchSegmentChanges, till, segmentName)]).then(flatMe => {
return Promise.all([json, greedyFetch(fetchSegmentChanges, till, segmentName, noCache)]).then(flatMe => {
return [flatMe[0], ...flatMe[1]];
});
}
Expand All @@ -34,11 +34,12 @@ export default function segmentChangesFetcherFactory(fetchSegmentChanges: IFetch
return function segmentChangesFetcher(
since: number,
segmentName: string,
noCache?: boolean,
// Optional decorator for `fetchMySegments` promise, such as timeout or time tracker
decorator?: (promise: Promise<ISegmentChangesResponse[]>) => Promise<ISegmentChangesResponse[]>
): Promise<ISegmentChangesResponse[]> {

let segmentsPromise = greedyFetch(fetchSegmentChanges, since, segmentName);
let segmentsPromise = greedyFetch(fetchSegmentChanges, since, segmentName, noCache);
if (decorator) segmentsPromise = decorator(segmentsPromise);

return segmentsPromise;
Expand Down
3 changes: 2 additions & 1 deletion src/sync/polling/fetchers/splitChangesFetcher.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,11 +10,12 @@ export default function splitChangesFetcherFactory(fetchSplitChanges: IFetchSpli

return function splitChangesFetcher(
since: number,
noCache?: boolean,
// Optional decorator for `fetchSplitChanges` promise, such as timeout or time tracker
decorator?: (promise: Promise<Response>) => Promise<Response>
) {

let splitsPromise = fetchSplitChanges(since);
let splitsPromise = fetchSplitChanges(since, noCache);
if (decorator) splitsPromise = decorator(splitsPromise);

return splitsPromise
Expand Down
12 changes: 9 additions & 3 deletions src/sync/polling/fetchers/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,12 +2,18 @@ import { ISplitChangesResponse, ISegmentChangesResponse } from '../../../dtos/ty

export type ISplitChangesFetcher = (
since: number,
decorator?: (promise: Promise<Response>) => Promise<Response>) => Promise<ISplitChangesResponse>
noCache?: boolean,
decorator?: (promise: Promise<Response>) => Promise<Response>
) => Promise<ISplitChangesResponse>

export type ISegmentChangesFetcher = (
since: number,
segmentName: string,
decorator?: (promise: Promise<ISegmentChangesResponse[]>) => Promise<ISegmentChangesResponse[]>) => Promise<ISegmentChangesResponse[]>
noCache?: boolean,
decorator?: (promise: Promise<ISegmentChangesResponse[]>) => Promise<ISegmentChangesResponse[]>
) => Promise<ISegmentChangesResponse[]>

export type IMySegmentsFetcher = (
decorator?: (promise: Promise<Response>) => Promise<Response>) => Promise<string[]>
noCache?: boolean,
decorator?: (promise: Promise<Response>) => Promise<Response>
) => Promise<string[]>
Original file line number Diff line number Diff line change
Expand Up @@ -69,7 +69,7 @@ test('splitChangesUpdater / factory', (done) => {
const readinessManager = readinessManagerFactory(EventEmitter);
const splitsEmitSpy = jest.spyOn(readinessManager.splits, 'emit');

const splitChangesUpdater = splitChangesUpdaterFactory(splitChangesFetcher, splitsCache, segmentsCache, readinessManager.splits);
const splitChangesUpdater = splitChangesUpdaterFactory(splitChangesFetcher, splitsCache, segmentsCache, readinessManager.splits, 1000, 1);

splitChangesUpdater().then((result) => {
expect(setChangeNumber.mock.calls.length).toBe(1);
Expand Down
16 changes: 9 additions & 7 deletions src/sync/polling/syncTasks/mySegmentsSyncTask.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,9 +10,10 @@ import { logFactory } from '../../../logger/sdkLogger';
import { IFetchMySegments } from '../../../services/types';
import mySegmentsFetcherFactory from '../fetchers/mySegmentsFetcher';
import { ISettings } from '../../../types';
import { SDK_SEGMENTS_ARRIVED } from '../../../readiness/constants';
const log = logFactory('splitio-sync:my-segments');

type IMySegmentsUpdater = (segmentList?: string[]) => Promise<boolean>
type IMySegmentsUpdater = (segmentList?: string[], noCache?: boolean) => Promise<boolean>

/**
* factory of MySegments updater (a.k.a, SegmentsSyncTask), a task that:
Expand Down Expand Up @@ -50,16 +51,16 @@ function mySegmentsUpdaterFactory(
// Notify update if required
if (splitsCache.usesSegments() && (shouldNotifyUpdate || readyOnAlreadyExistentState)) {
readyOnAlreadyExistentState = false;
segmentsEventEmitter.emit('SDK_SEGMENTS_ARRIVED');
segmentsEventEmitter.emit(SDK_SEGMENTS_ARRIVED);
}
}

function _mySegmentsUpdater(retry: number, segmentList?: string[]): Promise<boolean> {
function _mySegmentsUpdater(retry: number, segmentList?: string[], noCache?: boolean): Promise<boolean> {
const updaterPromise: Promise<boolean> = segmentList ?
// If segmentList is provided, there is no need to fetch mySegments
new Promise((res) => { updateSegments(segmentList); res(true); }) :
// If not provided, fetch mySegments
mySegmentsFetcher(_promiseDecorator).then(segments => {
mySegmentsFetcher(noCache, _promiseDecorator).then(segments => {
// Only when we have downloaded segments completely, we should not keep retrying anymore
startingUp = false;

Expand All @@ -74,7 +75,7 @@ function mySegmentsUpdaterFactory(
if (startingUp && retriesOnFailureBeforeReady > retry) {
retry += 1;
log.warn(`Retrying download of segments #${retry}. Reason: ${error}`);
return _mySegmentsUpdater(retry);
return _mySegmentsUpdater(retry); // no need to forward `segmentList` and `noCache` params
} else {
startingUp = false;
}
Expand All @@ -87,9 +88,10 @@ function mySegmentsUpdaterFactory(
* MySegments updater returns a promise that resolves with a `false` boolean value if it fails to fetch mySegments or synchronize them with the storage.
*
* @param {string[] | undefined} segmentList list of mySegments names to sync in the storage. If the list is `undefined`, it fetches them before syncing in the storage.
* @param {boolean | undefined} noCache true to revalidate data to fetch
*/
return function mySegmentsUpdater(segmentList?: string[]) {
return _mySegmentsUpdater(0, segmentList);
return function mySegmentsUpdater(segmentList?: string[], noCache?: boolean) {
return _mySegmentsUpdater(0, segmentList, noCache);
};

}
Expand Down
19 changes: 12 additions & 7 deletions src/sync/polling/syncTasks/segmentsSyncTask.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,10 +10,11 @@ import { logFactory } from '../../../logger/sdkLogger';
import segmentChangesFetcherFactory from '../fetchers/segmentChangesFetcher';
import { IFetchSegmentChanges } from '../../../services/types';
import { ISettings } from '../../../types';
import { SDK_SEGMENTS_ARRIVED } from '../../../readiness/constants';
const log = logFactory('splitio-sync:segment-changes');
const inputValidationLog = logFactory('', { displayAllErrors: true });

type ISegmentChangesUpdater = (segmentNames?: string[]) => Promise<boolean>
type ISegmentChangesUpdater = (segmentNames?: string[], noCache?: boolean, fetchOnlyNew?: boolean) => Promise<boolean>

/**
* factory of SegmentChanges updater (a.k.a, SegmentsSyncTask), a task that:
Expand Down Expand Up @@ -42,23 +43,27 @@ function segmentChangesUpdaterFactory(
* Thus, a false result doesn't imply that SDK_SEGMENTS_ARRIVED was not emitted.
*
* @param {string[] | undefined} segmentNames list of segment names to fetch. By passing `undefined` it fetches the list of segments registered at the storage
* @param {boolean | undefined} noCache true to revalidate data to fetch on a SEGMENT_UPDATE notifications.
* @param {boolean | undefined} fetchOnlyNew if true, only fetch the segments that not exists, i.e., which `changeNumber` is equal to -1.
* This param is used by SplitUpdateWorker on server-side SDK, to fetch new registered segments on SPLIT_UPDATE notifications.
*/
return function segmentChangesUpdater(segmentNames?: string[]) {
return function segmentChangesUpdater(segmentNames?: string[], noCache?: boolean, fetchOnlyNew?: boolean) {
log.debug('Started segments update');

// If not a segment name provided, read list of available segments names to be updated.
if (!segmentNames) segmentNames = segmentsCache.getRegisteredSegments();
let segments = segmentNames ? segmentNames : segmentsCache.getRegisteredSegments();
if (fetchOnlyNew) segments = segments.filter(segmentName => segmentsCache.getChangeNumber(segmentName) === -1);

// Async fetchers are collected here.
const updaters: Promise<number>[] = [];

for (let index = 0; index < segmentNames.length; index++) {
const segmentName = segmentNames[index];
for (let index = 0; index < segments.length; index++) {
const segmentName = segments[index];
const since = segmentsCache.getChangeNumber(segmentName);

log.debug(`Processing segment ${segmentName}`);

updaters.push(segmentChangesFetcher(since, segmentName, _promiseDecorator).then(function (changes) {
updaters.push(segmentChangesFetcher(since, segmentName, noCache, _promiseDecorator).then(function (changes) {
let changeNumber = -1;
changes.forEach(x => {
if (x.added.length > 0) segmentsCache.addToSegment(segmentName, x.added);
Expand All @@ -80,7 +85,7 @@ function segmentChangesUpdaterFactory(
// if at least one segment fetch successes, mark segments ready
if (findIndex(shouldUpdateFlags, v => v !== -1) !== -1 || readyOnAlreadyExistentState) {
readyOnAlreadyExistentState = false;
readiness.segments.emit('SDK_SEGMENTS_ARRIVED');
readiness.segments.emit(SDK_SEGMENTS_ARRIVED);
}
// if at least one segment fetch fails, return false to indicate that there was some error (e.g., invalid json, HTTP error, etc)
if (shouldUpdateFlags.indexOf(-1) !== -1) return false;
Expand Down
Loading