Skip to content
Draft
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
5 changes: 5 additions & 0 deletions .changeset/skip-clone-no-dedup.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
'@apollo/datasource-rest': patch
---

Add `shouldCloneParsedBodyForDeduplication()` so subclasses can skip per-caller cloning under high-fanout request deduplication. The default still clones concurrent consumers for isolation, and still skips the clone when a request-lifetime GET has only one consumer.
1 change: 1 addition & 0 deletions cspell-dict.txt
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ datasource
direnv
Fakeable
falsey
fanout
httpcache
instanceof
isplainobject
Expand Down
80 changes: 69 additions & 11 deletions src/RESTDataSource.ts
Original file line number Diff line number Diff line change
Expand Up @@ -195,6 +195,10 @@ export type RequestDeduplicationPolicy =
export abstract class RESTDataSource<CO extends CacheOptions = CacheOptions> {
protected httpCache: HTTPCache<CO>;
protected deduplicationPromises = new Map<string, Promise<any>>();
private deduplicationConsumers = new WeakMap<
Promise<any>,
{ count: number }
>();
baseURL?: string;
logger: Logger;

Expand Down Expand Up @@ -324,11 +328,14 @@ export abstract class RESTDataSource<CO extends CacheOptions = CacheOptions> {
'requestDeduplication'
>,
requestDeduplicationResult: RequestDeduplicationResult,
cloneParsedBody = true,
): DataSourceFetchResult<TResult> {
return {
...dataSourceFetchResult,
requestDeduplication: requestDeduplicationResult,
parsedBody: this.cloneParsedBody(dataSourceFetchResult.parsedBody),
parsedBody: cloneParsedBody
? this.cloneParsedBody(dataSourceFetchResult.parsedBody)
: dataSourceFetchResult.parsedBody,
};
}

Expand All @@ -337,6 +344,27 @@ export abstract class RESTDataSource<CO extends CacheOptions = CacheOptions> {
return cloneDeep(parsedBody);
}

/**
* Whether a caller of a deduplicated request should receive an independent
* clone of the parsed body.
*
* The default is `true`, which preserves isolation: mutating one caller's
* `parsedBody` does not affect other callers sharing the same HTTP request.
*
* Under high fan-out (many resolvers sharing one GET), that isolation costs
* one deep clone per consumer. Override and return `false` when callers treat
* the body as immutable to eliminate those clones. The default still skips
* cloning when `deduplicate-during-request-lifetime` has only a single
* consumer, because there is then nothing to isolate.
*
* Overriding `cloneParsedBody` to return its argument also skips these copies;
* this hook is the more explicit opt-out for "do not isolate deduplicated
* results."
*/
protected shouldCloneParsedBodyForDeduplication(): boolean {
return true;
}

protected shouldJSONSerializeBody(
body: RequestWithBody<CO>['body'],
): boolean {
Expand Down Expand Up @@ -583,33 +611,63 @@ export abstract class RESTDataSource<CO extends CacheOptions = CacheOptions> {
const previousRequestPromise = this.deduplicationPromises.get(
policy.deduplicationKey,
);
if (previousRequestPromise)
return previousRequestPromise.then((result) =>
this.cloneDataSourceFetchResult(result, {
if (previousRequestPromise) {
const previousConsumers = this.deduplicationConsumers.get(
previousRequestPromise,
);
if (previousConsumers) {
previousConsumers.count += 1;
}
return previousRequestPromise.then((result) => {
const requestDeduplicationResult: RequestDeduplicationResult = {
policy,
deduplicatedAgainstPreviousRequest: true,
}),
);
};
return this.cloneDataSourceFetchResult(
result,
requestDeduplicationResult,
this.shouldCloneParsedBodyForDeduplication(),
);
});
}

const thisRequestPromise = performRequest();
this.deduplicationPromises.set(
policy.deduplicationKey,
thisRequestPromise,
);
const consumers = { count: 1 };
this.deduplicationConsumers.set(thisRequestPromise, consumers);
try {
// The request promise needs to be awaited here rather than just
// returned. This ensures that the request completes before it's removed
// from the cache. Additionally, the use of finally here guarantees the
// deduplication cache is cleared in the event of an error during the
// request.
//
// Note: we could try to get fancy and only clone if no de-duplication
// happened (and we're "deduplicate-during-request-lifetime") but we
// haven't quite bothered yet.
return this.cloneDataSourceFetchResult(await thisRequestPromise, {
// Skip cloning when this is the only consumer and the result will not
// stay cached (`deduplicate-during-request-lifetime`). Concurrent
// waiters and `deduplicate-until-invalidated` still need a clone so
// callers can mutate independently.
const result = await thisRequestPromise;
const requestDeduplicationResult: RequestDeduplicationResult = {
policy,
deduplicatedAgainstPreviousRequest: false,
});
};
// Concurrent waiters and `deduplicate-until-invalidated` need isolation
// by default. A sole `deduplicate-during-request-lifetime` consumer
// does not. Subclasses can skip remaining clones via
// `shouldCloneParsedBodyForDeduplication`.
const needsIsolation =
consumers.count > 1 ||
policy.policy === 'deduplicate-until-invalidated';
const shouldCloneParsedBody =
needsIsolation && this.shouldCloneParsedBodyForDeduplication();
return this.cloneDataSourceFetchResult(
result,
requestDeduplicationResult,
shouldCloneParsedBody,
);
} finally {
if (policy.policy === 'deduplicate-during-request-lifetime') {
this.deduplicationPromises.delete(policy.deduplicationKey);
Expand Down
183 changes: 183 additions & 0 deletions src/__tests__/RESTDataSource.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1105,6 +1105,189 @@ describe('RESTDataSource', () => {
`);
});

it('does not clone a request-lifetime result with a single consumer', async () => {
let cloneCount = 0;
const dataSource = new (class extends RESTDataSource {
override baseURL = 'https://api.example.com';

protected override cloneParsedBody<TResult>(parsedBody: TResult) {
cloneCount += 1;
return super.cloneParsedBody(parsedBody);
}

getFoo(id: number) {
return this.fetch(`foo/${id}`);
}
})();

nock(apiUrl).get('/foo/1').reply(200, { hi: 42 });
nock(apiUrl).get('/foo/1').reply(200, { hi: 43 });

await dataSource.getFoo(1);
await dataSource.getFoo(1);
expect(cloneCount).toEqual(0);
});

it('clones parsed body for each concurrent consumer', async () => {
let cloneCount = 0;
const dataSource = new (class extends RESTDataSource {
override baseURL = 'https://api.example.com';

protected override cloneParsedBody<TResult>(parsedBody: TResult) {
cloneCount += 1;
return super.cloneParsedBody(parsedBody);
}

getFoo(id: number) {
return this.fetch(`foo/${id}`);
}
})();

nock(apiUrl).get('/foo/1').reply(200, { hi: 42 });

await Promise.all([dataSource.getFoo(1), dataSource.getFoo(1)]);
expect(cloneCount).toEqual(2);
});

it('clones parsed body once per concurrent consumer under high fan-out', async () => {
const consumerCount = 100;
let cloneCount = 0;
const dataSource = new (class extends RESTDataSource {
override baseURL = 'https://api.example.com';

protected override cloneParsedBody<TResult>(parsedBody: TResult) {
cloneCount += 1;
return super.cloneParsedBody(parsedBody);
}

getFoo(id: number) {
return this.fetch(`foo/${id}`);
}
})();

nock(apiUrl).get('/foo/1').reply(200, { hi: 42 });

await Promise.all(
Array.from({ length: consumerCount }, () => dataSource.getFoo(1)),
);
expect(cloneCount).toEqual(consumerCount);
});

it('does not clone under high fan-out when shouldCloneParsedBodyForDeduplication is false', async () => {
const consumerCount = 100;
let cloneCount = 0;
const dataSource = new (class extends RESTDataSource {
override baseURL = 'https://api.example.com';

protected override shouldCloneParsedBodyForDeduplication(): boolean {
return false;
}

protected override cloneParsedBody<TResult>(parsedBody: TResult) {
cloneCount += 1;
return super.cloneParsedBody(parsedBody);
}

getFoo(id: number) {
return this.fetch<{ hi: number }>(`foo/${id}`);
}
})();

nock(apiUrl).get('/foo/1').reply(200, { hi: 42 });

const results = await Promise.all(
Array.from({ length: consumerCount }, () => dataSource.getFoo(1)),
);
expect(cloneCount).toEqual(0);
expect(results).toHaveLength(consumerCount);
for (const result of results) {
expect(result.parsedBody).toEqual({ hi: 42 });
expect(result.parsedBody).toBe(results[0].parsedBody);
}
});

it('clones parsed body for deduplicate-until-invalidated so later callers are isolated', async () => {
let cloneCount = 0;
const dataSource = new (class extends RESTDataSource {
override baseURL = 'https://api.example.com';
protected override requestDeduplicationPolicyFor(
url: URL,
request: WithRequired<RequestOptions, 'method'>,
): RequestDeduplicationPolicy {
const p = super.requestDeduplicationPolicyFor(url, request);
return p.policy === 'deduplicate-during-request-lifetime'
? {
policy: 'deduplicate-until-invalidated',
deduplicationKey: p.deduplicationKey,
}
: p;
}

protected override cloneParsedBody<TResult>(parsedBody: TResult) {
cloneCount += 1;
return super.cloneParsedBody(parsedBody);
}

getFoo(id: number) {
return this.fetch<{ hi: number }>(`foo/${id}`);
}
})();

nock(apiUrl).get('/foo/1').reply(200, { hi: 42 });

const r1 = await dataSource.getFoo(1);
r1.parsedBody.hi = 99;
const r2 = await dataSource.getFoo(1);
expect(r1.parsedBody.hi).toEqual(99);
expect(r2.parsedBody.hi).toEqual(42);
expect(cloneCount).toEqual(2);
});

// Opting out of cloning under `deduplicate-until-invalidated` is the
// most dangerous mode: the shared parsed result stays cached and later
// callers receive the same object. Mutations are intentionally visible
// across sequential callers (aliased, not isolated).
it('does not clone under deduplicate-until-invalidated when shouldCloneParsedBodyForDeduplication is false', async () => {
let cloneCount = 0;
const dataSource = new (class extends RESTDataSource {
override baseURL = 'https://api.example.com';
protected override requestDeduplicationPolicyFor(
url: URL,
request: WithRequired<RequestOptions, 'method'>,
): RequestDeduplicationPolicy {
const p = super.requestDeduplicationPolicyFor(url, request);
return p.policy === 'deduplicate-during-request-lifetime'
? {
policy: 'deduplicate-until-invalidated',
deduplicationKey: p.deduplicationKey,
}
: p;
}

protected override shouldCloneParsedBodyForDeduplication(): boolean {
return false;
}

protected override cloneParsedBody<TResult>(parsedBody: TResult) {
cloneCount += 1;
return super.cloneParsedBody(parsedBody);
}

getFoo(id: number) {
return this.fetch<{ hi: number }>(`foo/${id}`);
}
})();

nock(apiUrl).get('/foo/1').reply(200, { hi: 42 });

const r1 = await dataSource.getFoo(1);
r1.parsedBody.hi = 99;
const r2 = await dataSource.getFoo(1);
expect(cloneCount).toEqual(0);
expect(r2.parsedBody).toBe(r1.parsedBody);
expect(r2.parsedBody.hi).toEqual(99);
});

it('does not deduplicate non-GET requests by default', async () => {
const dataSource = new (class extends RESTDataSource {
override baseURL = 'https://api.example.com';
Expand Down