Repository navigation
Conversation
…Actor.childRuns()
|
See more at https://github.com/apify/apify-sdk-js/actions/runs/37911125694#summary-113756252850 |
barjin
left a comment
There was a problem hiding this comment.
A few ideas / comments where this implementation differs from the Orchestrator ⬇️
| const tracked = await this.#childRunTracker.get(runName); | ||
| const trackedRun = tracked ? await client.run(tracked.runId).get() : undefined; | ||
|
|
||
| if (trackedRun && ['SUCCEEDED', 'READY', 'RUNNING'].includes(trackedRun.status)) { |
There was a problem hiding this comment.
Every other run status (aborting / timed-out) is currently replaced by a completely new run instance.
Should we attempt to resurrect these?
There was a problem hiding this comment.
Bonus question for resurrected runs - shall we try implementing the 'inherit' timeout mode (pass the remaining timeout to the resurrected child run)?
There's no Actor.resurrect() method in SDK now (is this intentional?)
There was a problem hiding this comment.
Every other run status (aborting / timed-out) is currently replaced by a completely new run instance. Should we attempt to resurrect these?
Yes, I'd rather resurrect them.
Bonus question for resurrected runs - shall we try implementing the 'inherit' timeout mode (pass the remaining timeout to the resurrected child run)?
It makes sense. Let's try to implement it.
There's no Actor.resurrect() method in SDK now (is this intentional?)
I'm not aware of this being intentionally missing, so I'd add it and use it (Python SDK is missing it as well).
Co-authored-by: Martin Adámek <banan23@gmail.com>
Co-authored-by: Martin Adámek <banan23@gmail.com>
Co-authored-by: Martin Adámek <banan23@gmail.com>
`Date` and `URL` inputs were hashed as `{}` and function inputs were dropped, so inputs differing only in those
values got the same checksum and the earlier run was resumed. Hash the JSON form `apify-client` sends instead,
and sort keys by code units, as `localeCompare` depends on the locale.
|
Lets compare this with apify/apify-sdk-python#1149 |
|
Compared with the Python counterpart, which is split into apify/apify-sdk-python#1149 (named runs) and apify/apify-sdk-python#1151 (
Things we should settle across both SDKs:
|
|
Python PRs (apify/apify-sdk-python#1149, apify/apify-sdk-python#1151) have been changed to match JS behavior in two places:
I'm still thinking about the rest. |
janbuchar
left a comment
There was a problem hiding this comment.
I have a bunch of code cleanup suggestions, but nothing critical. The PR looks fine to me.
bc1cc0b to
c83bcc6
Compare
|
|
||
| return this.locks.get(runName)!.runExclusive(async () => { | ||
| const tracked = await this.get(runName); | ||
| const trackedRun = tracked ? await client.run(tracked.runId).get() : undefined; |
There was a problem hiding this comment.
A 404 here right after a start (e.g. a concurrent call, or a migration moments after start()) can just be a lagging API replica, and we would start a duplicate run. Python retries get() for 3 s before treating the run as LOST.
See apify/apify-sdk-python#1149 (comment).
✍️ Drafted by Claude Code
| const tracked = await this.get(runName); | ||
| const trackedRun = tracked ? await client.run(tracked.runId).get() : undefined; | ||
|
|
||
| if (trackedRun && ['SUCCEEDED', 'READY', 'RUNNING'].includes(trackedRun.status)) { |
There was a problem hiding this comment.
An ABORTING / TIMING-OUT run falls through to start(), so the old and new runs overlap for a while. Python waits for such a run to finish (waitForFinish()) and then handles it as ABORTED / TIMED-OUT.
✍️ Drafted by Claude Code
| const trackedRun = tracked ? await client.run(tracked.runId).get() : undefined; | ||
|
|
||
| if (trackedRun && ['SUCCEEDED', 'READY', 'RUNNING'].includes(trackedRun.status)) { | ||
| await this.verifyRequest(runName, request); |
There was a problem hiding this comment.
The checksum is only checked on this reattach path, so a name whose recorded run failed or got lost accepts a different input, and the new checksum overwrites the old one. Python raises in that case too, so a name stays bound to one Actor/task + input. Which way we want to go?
✍️ Drafted by Claude Code
| return { run: trackedRun, resumed: true }; | ||
| } | ||
|
|
||
| const run = await start(); |
There was a problem hiding this comment.
ABORTED / TIMED-OUT runs end up here and get replaced by a new run, but we agreed to resurrect them. Python resurrects them with the call's build, memory, timeout and max charge. JS needs an Actor.resurrect() for that first.
This can be addressed once this #769 is merged.
✍️ Drafted by Claude Code
|
Compared with the Python counterpart (apify/apify-sdk-python#1149 for named runs, apify/apify-sdk-python#1151 for Full comparison table:
✍️ Drafted by Claude Code |
`Actor.start`, `Actor.call`, `Actor.start_task` and `Actor.call_task` accept an optional `run_name`. Right after the child starts, the SDK records it in the default key-value store under `__ACTOR_CHILD_RUNS`. After a migration or resurrection of the parent, the same call finds the recorded run and reuses it: - `READY` / `RUNNING`: reattaches to it - `SUCCEEDED`: returns it as is - `ABORTED` / `TIMED-OUT`: resurrects it with the call's build, memory, timeout, `max_items` and max charge (an `ABORTING` / `TIMING-OUT` run is waited for first) - `FAILED`, or the run no longer exists: starts a new run and moves the old one to `history` The record has the same shape as in the JS SDK (apify/apify-sdk-js#759): `runId`, `status`, `startedAt`, `checksum` and `history`. The checksum covers the Actor or task ID and the input, computed the same way as in JS. Reusing a name for a different Actor, task or input raises `ValueError`. Concurrent calls under one name in the same process share a single run. A named `call` on a reattached or resurrected run streams only new log lines. If the API reports a recorded run as missing, the SDK retries for 3 seconds before replacing it, because a run started moments ago may not be on every API replica yet. A hard kill between the platform starting the child and the registry write can still orphan the child. Closing that gap needs an idempotency key on the run-start endpoint. Closes: #1127 *✍️ Drafted by Claude Code*
`Actor.child_runs` is a sync property that maps each run name to a `RunClientAsync` for the current run under that name. It has the same shape as `Actor.childRuns` in the JS SDK (apify/apify-sdk-js#759). The client covers the run and its storages: ```python run = await Actor.child_runs['scrape-eu'].wait_for_finish() ``` The SDK reads the registry from the default key-value store when the Actor initializes, so runs recorded before a migration or resurrection are included and reading the property makes no API calls. That adds one record read to every init. A malformed record fails the init and leaves the Actor uninitialized. Runs started without a `run_name` aren't tracked. Runs that were replaced under a name stay in the `history` of the `__ACTOR_CHILD_RUNS` record and aren't exposed. A run started or reattached with a custom `token` gets a client with that token. A run recorded before a migration and not reattached since gets the default client, because the token isn't stored. It also fixes a duplicate `max_items` deprecation warning from #1149: a named start that resurrects its recorded run no longer warns a second time from inside the SDK. Docs are in #1160. Closes: #1129 *✍️ Drafted by Claude Code*
| const { waitSecs, log, ...startOptions } = rest; | ||
| const { run, resumed } = await this.#childRunTracker.startOrResumeChildRun( | ||
| client, | ||
| runName, | ||
| { type: 'actor', id: actorId, input }, | ||
| async () => client.actor(actorId).start(input, { ...startOptions, runTimeoutSecs }), | ||
| ); | ||
|
|
||
| // The earlier part of a resumed run's log was already redirected before the migration. | ||
| const streamedLog = await client.run(run.id).getStreamedLog({ toLog: log, fromStart: !resumed }); | ||
| streamedLog?.start(); | ||
| try { | ||
| const finishedRun = await client.run(run.id).waitForFinish({ waitSecs }); |
There was a problem hiding this comment.
[medium] With runName, signal and timeoutSecs only reach start(). The unnamed path (client call()) also passes them to waitForFinish / getStreamedLog and defaults timeoutSecs: 'noTimeout'. Aborting the signal doesn't stop Actor.call(..., { runName, signal }); it hangs until the child finishes. And start() falls back to the 'medium' request timeout, which can cut off a start with waitForResources. Same in callTask() below. Fix: forward { waitSecs, timeoutSecs, signal } to waitForFinish, signal to getStreamedLog, and default timeoutSecs like the client does.
Repro (with runName times out, without passes)
import { createServer } from 'node:http';
import type { AddressInfo } from 'node:net';
import { MemoryStorageBackend } from '@crawlee/core';
import { Configuration } from 'apify';
import type { Run as ActorRun } from 'apify-client';
import { ActorClient } from 'apify-client';
import { createIsolatedActor } from '../createIsolatedActor.js';
const server = createServer(() => {});
let apiBaseUrl: string;
beforeAll(async () => {
await new Promise<void>((resolve) => server.listen(0, '127.0.0.1', resolve));
apiBaseUrl = `http://127.0.0.1:${(server.address() as AddressInfo).port}/`;
});
afterAll(() => { server.closeAllConnections(); server.close(); });
const run = { id: 'child-run', status: 'RUNNING', startedAt: new Date() } as ActorRun;
test.each([['without runName', {}], ['with runName', { runName: 'child' }]])(
'call() %s rejects once its abort signal fires',
async (_label, extra) => {
vitest.spyOn(ActorClient.prototype, 'start').mockResolvedValue(run);
const { actor } = createIsolatedActor({
config: new Configuration({ apiBaseUrl, token: 'token' }),
storageClient: new MemoryStorageBackend(),
});
const signal = AbortSignal.timeout(200);
await expect(actor.call('some-actor', {}, { signal, log: null, ...extra })).rejects.toThrow();
},
3000,
);| const { waitSecs, ...startOptions } = rest; | ||
| const { run } = await this.#childRunTracker.startOrResumeChildRun( | ||
| client, | ||
| runName, | ||
| { type: 'task', id: taskId, input }, | ||
| async () => client.task(taskId).start(input, { ...startOptions, runTimeoutSecs }), | ||
| ); | ||
| const finishedRun = await client.run(run.id).waitForFinish({ waitSecs }); |
There was a problem hiding this comment.
[medium] Same as in call(): signal / timeoutSecs never reach waitForFinish.
| }); | ||
| log.debug(`Default storages purged`); | ||
|
|
||
| await this.#childRunTracker.load(); |
There was a problem hiding this comment.
[low] This adds a getRecord('__ACTOR_CHILD_RUNS') (plus a store get) to every Actor.init() on the platform, also for Actors that never use runName, and init fails if the read fails. Is that cost intended for the sync childRuns getter? Alternative: load lazily on the first named start/call/callTask.
| * | ||
| * If a `READY`, `RUNNING` or `SUCCEEDED` run is already tracked under this name, it is returned instead of | ||
| * starting a new one, including after a migration or resurrection of this run. A run that failed, was aborted | ||
| * or timed out is replaced by a new one. Reusing the name for a different Actor, task or input throws. |
There was a problem hiding this comment.
[low] It only throws while the tracked run is returned (READY/RUNNING/SUCCEEDED). Otherwise it's silently replaced, which the "allowed when the tracked run is replaced" test pins down. Suggest "...throws while the tracked run is returned."
| static #instance?: Actor; | ||
|
|
||
| /** | ||
| * Tracks child runs of this Actor instance. |
There was a problem hiding this comment.
[nit] Restates the field name, can be dropped.
| } | ||
|
|
||
| /** | ||
| * A FIFO mutex that a critical section may re-enter from a nested call — the charge lock is taken by |
There was a problem hiding this comment.
[nit] Doc still describes the charge lock / pushData(), but it's a shared util now (also per-name child-run locks), so make it general. Also the node:async_hooks import on L18 should go with the other builtins.
| ); | ||
| } | ||
|
|
||
| private async persist(trackedRuns: Record<string, TrackedChildRun>) { |
There was a problem hiding this comment.
[nit] persist() always gets this.loadedRuns, so the param is redundant. trackedRuns (promise) and loadedRuns also hold the same state, could be one field.
| import log from '@apify/log'; | ||
|
|
||
| import type { TrackedChildRun } from '../../src/child_run_tracker.js'; | ||
| import { createIsolatedActor } from '../createIsolatedActor'; |
There was a problem hiding this comment.
Holy guacamole, can we set up a lint rule for this?
There was a problem hiding this comment.
I've set up a hook for myself months ago, so it didn't bother me, but sure, we can: #770
Adds a
runNameoption toActor.start,callandcallTaskthat reattaches to the recorded child run after a restart, andActor.childRuns()to read them.Closes #738, closes #740