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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 17 additions & 0 deletions docs/02_concepts/06_timeouts.md
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,23 @@ A number above `timeoutMaxSecs` is capped at it, and the client logs a warning.

Methods that start an Actor run keep the API's run timeout apart from the request timeout. The `runTimeoutSecs` option of `start()`, `call()` and `resurrect()` bounds how long the run may execute on the platform, while `timeoutSecs` bounds the request that starts it.

## Aborting a call

Every method that accepts `timeoutSecs` also accepts a `signal` option that takes an `AbortSignal`. Once the signal aborts, the client ends the request in flight and skips any remaining retries. It sends no further poll, and the method rejects with the signal's `reason`. Use it to stop a call from a shutdown handler, or to give a polling method such as `call()` an overall deadline:

```js
const controller = new AbortController();
process.once('SIGTERM', () => controller.abort());

// Stops waiting for the run once the process receives SIGTERM. The run itself keeps going on the platform.
const run = await client.actor('my-actor').call(input, { signal: controller.signal });

// Gives up on the whole wait after 10 minutes, however many polls it takes.
const finishedRun = await client.run('my-run-id').waitForFinish({ signal: AbortSignal.timeout(600_000) });
```

Aborting the signal only stops the client. To stop an Actor run on the platform, call <ApiLink to="class/RunClient#abort">`RunClient.abort()`</ApiLink>.

## Interaction with retries

Timeouts work together with [retries](./02_error-handling.md#retries-with-exponential-backoff). A request that times out counts as a failed attempt and is retried, up to `maxRetries` times. The timeout applies to each attempt on its own, and doubles with every retry up to `timeoutMaxSecs`, so a request that timed out at 5 seconds gets 10 seconds on the second attempt and 20 on the third.
14 changes: 8 additions & 6 deletions docs/public-api/apify-client.api.md
Original file line number Diff line number Diff line change
Expand Up @@ -3599,11 +3599,11 @@ export interface RequestQueueUserOptions {
// Not exported by the entry point; reachable only as a referenced type.
// @public
class ResourceClient extends ApiClient {
protected deleteResource(timeoutSecs: Timeout): Promise<void>;
protected getResource<T, R>(schema: z.ZodType, options: T, timeoutSecs: Timeout): Promise<R | undefined>;
protected deleteResource(timeoutSecs: Timeout, signal?: AbortSignal): Promise<void>;
protected getResource<T, R>(schema: z.ZodType, options: T, timeoutSecs: Timeout, signal?: AbortSignal): Promise<R | undefined>;
protected timeoutForWaitForFinish(timeoutSecs: Timeout | undefined, tier: TimeoutTier, waitForFinishSecs: number | undefined): Timeout;
// (undocumented)
protected updateResource<T, R>(schema: z.ZodType, newFields: T, timeoutSecs: Timeout): Promise<R>;
protected updateResource<T, R>(schema: z.ZodType, newFields: T, timeoutSecs: Timeout, signal?: AbortSignal): Promise<R>;
protected waitForJobFinish<R extends {
status: (typeof ACT_JOB_STATUSES)[keyof typeof ACT_JOB_STATUSES];
}>(schema: z.ZodType, options?: WaitForFinishOptions): Promise<R>;
Expand All @@ -3613,11 +3613,11 @@ class ResourceClient extends ApiClient {
// @public
class ResourceCollectionClient extends ApiClient {
// (undocumented)
protected createResource<D, R>(schema: z.ZodType, resource: D, timeoutSecs: Timeout): Promise<R>;
protected createResource<D, R>(schema: z.ZodType, resource: D, timeoutSecs: Timeout, signal?: AbortSignal): Promise<R>;
// (undocumented)
protected getOrCreateResource<D, R>(schema: z.ZodType, name: string | undefined, resource: D | undefined, timeoutSecs: Timeout): Promise<R>;
protected getOrCreateResource<D, R>(schema: z.ZodType, name: string | undefined, resource: D | undefined, timeoutSecs: Timeout, signal?: AbortSignal): Promise<R>;
// (undocumented)
protected listResources<T, R>(schema: z.ZodType, options: T | undefined, timeoutSecs: Timeout): Promise<R>;
protected listResources<T, R>(schema: z.ZodType, options: T | undefined, timeoutSecs: Timeout, signal?: AbortSignal): Promise<R>;
protected listResourcesPaginated<T extends PaginationOptions & TimeoutOptions, Data, R extends PaginatedResponse<Data>>(schema: z.ZodType, options: T, defaultTimeoutSecs: Timeout): AsyncIterable<Data> & Promise<R>;
}

Expand Down Expand Up @@ -3882,6 +3882,7 @@ export class StreamedLog {
export interface StreamedLogOptions {
fromStart?: boolean;
logClient: LogClient;
signal?: AbortSignal;
toLog: Log;
}

Expand Down Expand Up @@ -4020,6 +4021,7 @@ export type Timeout = TimeoutTier | 'noTimeout' | number;

// @public
export interface TimeoutOptions {
signal?: AbortSignal;
timeoutSecs?: Timeout;
}

Expand Down
4 changes: 2 additions & 2 deletions src/base/api_client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -104,8 +104,8 @@ export abstract class ApiClient {
};

// `chunkSize` only sizes this loop's requests; it is not an API parameter, so it must not reach
// `buildParams()` and the query string. The same goes for `timeoutSecs`, which callers take out before
// calling this, since it also picks the timeout of every page request.
// `buildParams()` and the query string. The same goes for `timeoutSecs` and `signal`, which callers take out
// before calling this, since they also apply to every page request.
const { chunkSize, ...listOptions } = options;

const paginatedListPromise = getPaginatedList({
Expand Down
32 changes: 23 additions & 9 deletions src/base/resource_client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@ import type { z } from 'zod';
import type { ApifyApiError } from '../apify_api_error.js';
import type { ApifyRequestConfig } from '../http_clients/index.js';
import type { Timeout, TimeoutOptions, TimeoutTier } from '../timeouts.js';
import { catchNotFoundForResourceOrThrow, catchNotFoundOrThrow, parseResponse } from '../utils.js';
import { catchNotFoundForResourceOrThrow, catchNotFoundOrThrow, parseResponse, sleep } from '../utils.js';
import { ApiClient } from './api_client.js';

/**
Expand Down Expand Up @@ -47,12 +47,18 @@ export class ResourceClient extends ApiClient {
* A 404 resolves to `undefined` only when the client names its resource by ID. A chained client without one, such
* as `run.dataset()`, throws it instead (see `catchNotFoundForResourceOrThrow()`).
*/
protected async getResource<T, R>(schema: z.ZodType, options: T, timeoutSecs: Timeout): Promise<R | undefined> {
protected async getResource<T, R>(
schema: z.ZodType,
options: T,
timeoutSecs: Timeout,
signal?: AbortSignal,
): Promise<R | undefined> {
const requestOpts: ApifyRequestConfig = {
url: this.buildUrl(),
method: 'GET',
params: this.buildParams(options),
timeoutSecs,
signal,
};
try {
const response = await this.httpClient.call(requestOpts);
Expand All @@ -64,13 +70,19 @@ export class ResourceClient extends ApiClient {
return undefined;
}

protected async updateResource<T, R>(schema: z.ZodType, newFields: T, timeoutSecs: Timeout): Promise<R> {
protected async updateResource<T, R>(
schema: z.ZodType,
newFields: T,
timeoutSecs: Timeout,
signal?: AbortSignal,
): Promise<R> {
const response = await this.httpClient.call({
url: this.buildUrl(),
method: 'PUT',
params: this.buildParams(),
data: newFields,
timeoutSecs,
signal,
});
return parseResponse<R>(response, schema);
}
Expand All @@ -79,13 +91,14 @@ export class ResourceClient extends ApiClient {
* A 404 is swallowed, keeping the DELETE idempotent, only when the client names its resource by ID. A chained client
* without one throws it instead (see `catchNotFoundForResourceOrThrow()`).
*/
protected async deleteResource(timeoutSecs: Timeout): Promise<void> {
protected async deleteResource(timeoutSecs: Timeout, signal?: AbortSignal): Promise<void> {
try {
await this.httpClient.call({
url: this.buildUrl(),
method: 'DELETE',
params: this.buildParams(),
timeoutSecs,
signal,
});
} catch (err) {
catchNotFoundForResourceOrThrow(err as ApifyApiError, this.id);
Expand All @@ -100,7 +113,7 @@ export class ResourceClient extends ApiClient {
schema: z.ZodType,
options: WaitForFinishOptions = {},
): Promise<R> {
const { waitSecs = MAX_WAIT_FOR_FINISH, timeoutSecs = 'noTimeout' } = options;
const { waitSecs = MAX_WAIT_FOR_FINISH, timeoutSecs = 'noTimeout', signal } = options;
const waitMillis = waitSecs * 1000;
let job: R | undefined;

Expand All @@ -123,6 +136,7 @@ export class ResourceClient extends ApiClient {
method: 'GET',
params: this.buildParams({ waitForFinish }),
timeoutSecs,
signal,
};
try {
const response = await this.httpClient.call(requestOpts);
Expand All @@ -134,10 +148,10 @@ export class ResourceClient extends ApiClient {

// It might take some time for database replicas to get up-to-date,
// so getRun() might return null. Wait a little bit and try it again.
if (!job)
await new Promise((resolve) => {
setTimeout(resolve, 250);
});
if (!job) {
await sleep(250, signal);
signal?.throwIfAborted();
}
} while (shouldRepeat());

if (!job) {
Expand Down
28 changes: 22 additions & 6 deletions src/base/resource_collection_client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -13,41 +13,55 @@ export class ResourceCollectionClient extends ApiClient {
/**
* @private
*/
protected async listResources<T, R>(schema: z.ZodType, options: T | undefined, timeoutSecs: Timeout): Promise<R> {
protected async listResources<T, R>(
schema: z.ZodType,
options: T | undefined,
timeoutSecs: Timeout,
signal?: AbortSignal,
): Promise<R> {
const response = await this.httpClient.call({
url: this.buildUrl(),
method: 'GET',
params: this.buildParams(options),
timeoutSecs,
signal,
});
return parseResponse<R>(response, schema);
}

/**
* Returns async iterator to iterate through all items and Promise that can be awaited to get first page of results.
* `defaultTimeoutSecs` applies to every page request unless `options.timeoutSecs` overrides it.
* `defaultTimeoutSecs` applies to every page request unless `options.timeoutSecs` overrides it, and
* `options.signal` goes to every page request.
*/
protected listResourcesPaginated<
T extends PaginationOptions & TimeoutOptions,
Data,
R extends PaginatedResponse<Data>,
>(schema: z.ZodType, options: T, defaultTimeoutSecs: Timeout): AsyncIterable<Data> & Promise<R> {
// `timeoutSecs` only times the page requests; it is not an API parameter, so it must not reach the query string.
const { timeoutSecs = defaultTimeoutSecs, ...listOptions } = options;
// `timeoutSecs` and `signal` only apply to the page requests; they are not API parameters, so they must not
// reach the query string.
const { timeoutSecs = defaultTimeoutSecs, signal, ...listOptions } = options;

return this.listPaginatedFromCallback(
async (pageOptions?: T) => this.listResources<T, R>(schema, pageOptions, timeoutSecs),
async (pageOptions?: T) => this.listResources<T, R>(schema, pageOptions, timeoutSecs, signal),
listOptions as T,
);
}

protected async createResource<D, R>(schema: z.ZodType, resource: D, timeoutSecs: Timeout): Promise<R> {
protected async createResource<D, R>(
schema: z.ZodType,
resource: D,
timeoutSecs: Timeout,
signal?: AbortSignal,
): Promise<R> {
const response = await this.httpClient.call({
url: this.buildUrl(),
method: 'POST',
params: this.buildParams(),
data: resource,
timeoutSecs,
signal,
});
return parseResponse<R>(response, schema);
}
Expand All @@ -57,13 +71,15 @@ export class ResourceCollectionClient extends ApiClient {
name: string | undefined,
resource: D | undefined,
timeoutSecs: Timeout,
signal?: AbortSignal,
): Promise<R> {
const response = await this.httpClient.call({
url: this.buildUrl(),
method: 'POST',
params: this.buildParams({ name }),
data: resource,
timeoutSecs,
signal,
});
return parseResponse<R>(response, schema);
}
Expand Down
22 changes: 1 addition & 21 deletions src/http_clients/base.ts
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@ import {
isStream,
MIN_COMPRESS_BYTES,
parseArgument,
sleep,
toBytes,
version,
} from '../utils.js';
Expand Down Expand Up @@ -677,27 +678,6 @@ export abstract class HttpClient {
}
}

/**
* Resolves after `millis`, or right away once `signal` aborts.
*/
async function sleep(millis: number, signal?: AbortSignal): Promise<void> {
return new Promise((resolve) => {
if (signal?.aborted) {
resolve();
return;
}
const onAbort = () => {
clearTimeout(timer);
resolve();
};
const timer = setTimeout(() => {
signal?.removeEventListener('abort', onAbort);
resolve();
}, millis);
signal?.addEventListener('abort', onAbort, { once: true });
});
}

/**
* Looks a header up by name, compared case-insensitively.
*/
Expand Down
Loading
Loading