import { APIFY_BUILD_SCHEMA, APIFY_JOB_STATUS_SCHEMA, APIFY_RUN_SCHEMA, mapApifyBuild, mapApifyRun, } from "@automate.ax/integration-contracts/apify" import * as z from "zod" import { defineAction } from "../../automation/actions" import { getApifyClient } from "./lib" const APIFY_RUN_OPTIONS_SCHEMA = z.object({ /** Actor build tag or number. */ build: z.string().min(1).optional(), /** Maximum number of paid dataset items. */ maxItems: z.number().int().positive().optional(), /** Maximum total run charge in US dollars. */ maxTotalChargeUsd: z.number().positive().optional(), /** Run memory in megabytes. */ memoryMbytes: z .number() .int() .min(128) .max(32_768) .refine(isPowerOfTwo, "Memory must be a power of two.") .optional(), /** Restart the run when its container exits with an error. */ restartOnError: z.boolean().optional(), /** Maximum run duration in seconds. Zero disables the timeout. */ timeoutSecs: z.number().int().nonnegative().optional(), }) const APIFY_RUN_LIST_OUTPUT_SCHEMA = z.object({ count: z.number().int().nonnegative(), desc: z.boolean(), items: APIFY_RUN_SCHEMA.array(), limit: z.number().int().nonnegative(), offset: z.number().int().nonnegative(), total: z.number().int().nonnegative(), }) const APIFY_DATASET_ITEMS_OUTPUT_SCHEMA = z.object({ count: z.number().int().nonnegative(), desc: z.boolean(), items: z.record(z.string(), z.json()).array(), limit: z.number().int().nonnegative(), offset: z.number().int().nonnegative(), total: z.number().int().nonnegative(), }) /** Starts an Apify Actor without waiting for it to finish. */ export const startApifyActor = defineAction("Start Apify Actor") .describe("Starts an Apify Actor and returns its run immediately.") .account("apify") .input( APIFY_RUN_OPTIONS_SCHEMA.extend({ /** Actor ID or username/name reference. */ actorId: z.string().min(1), /** Actor-defined JSON input. */ input: z.json().optional(), }), ) .output(APIFY_RUN_SCHEMA) .retry({ replaySafety: "unsafe" }) .handler(async ({ account, input }) => { const { actorId, input: actorInput, ...options } = input return mapApifyRun( await getApifyClient(account.secret) .actor(actorId) .start(actorInput, { ...mapRunOptions(options), waitForFinish: 0, }), ) }) /** Starts an Apify task without waiting for it to finish. */ export const startApifyTask = defineAction("Start Apify task") .describe("Starts a saved Apify task and returns its run immediately.") .account("apify") .input( APIFY_RUN_OPTIONS_SCHEMA.extend({ /** JSON properties that override the task's saved input. */ input: z.record(z.string(), z.json()).optional(), /** Task ID or username/name reference. */ taskId: z.string().min(1), }), ) .output(APIFY_RUN_SCHEMA) .retry({ replaySafety: "unsafe" }) .handler(async ({ account, input }) => { const { input: taskInput, taskId, ...options } = input return mapApifyRun( await getApifyClient(account.secret) .task(taskId) .start(taskInput, { ...mapRunOptions(options), waitForFinish: 0, }), ) }) /** Gets one Apify Actor run. */ export const getApifyRun = defineAction("Get Apify Actor run") .describe("Gets the current state and storage IDs for one Apify Actor run.") .account("apify") .input(z.object({ runId: z.string().min(1) })) .output(APIFY_RUN_SCHEMA) .retry({ replaySafety: "safe" }) .handler(async ({ account, input }) => mapApifyRun( await requireResource( getApifyClient(account.secret).run(input.runId).get(), "Actor run", input.runId, ), ), ) /** Lists Apify Actor runs across an account, Actor, or task. */ export const listApifyRuns = defineAction("List Apify Actor runs") .describe("Lists Apify Actor runs with status, time, and resource filters.") .account("apify") .input( z.object({ /** Limit results to one Actor ID or username/name reference. */ actorId: z.string().min(1).optional(), /** Return newest runs first. */ desc: z.boolean().optional(), /** Maximum number of runs returned by this page. */ limit: z.number().int().min(1).max(1_000).optional(), /** Number of runs to skip. */ offset: z.number().int().nonnegative().optional(), /** Return runs started after this ISO 8601 timestamp. */ startedAfter: z.iso.datetime().optional(), /** Return runs started before this ISO 8601 timestamp. */ startedBefore: z.iso.datetime().optional(), /** Filter by one or more run statuses. */ status: z .union([ APIFY_JOB_STATUS_SCHEMA, APIFY_JOB_STATUS_SCHEMA.array().min(1), ]) .optional(), /** Limit results to one saved task. */ taskId: z.string().min(1).optional(), }), ) .output(APIFY_RUN_LIST_OUTPUT_SCHEMA) .retry({ replaySafety: "safe" }) .handler(async ({ account, input }) => { const { actorId, taskId, ...options } = input if (actorId && taskId) { throw new Error("List Apify runs by either actorId or taskId, not both.") } const client = getApifyClient(account.secret) // Select the scoped collection once before applying shared pagination. const runs = actorId ? client.actor(actorId).runs() : taskId ? client.task(taskId).runs() : client.runs() const result = await runs.list(options) return APIFY_RUN_LIST_OUTPUT_SCHEMA.parse({ ...result, items: result.items.map(mapApifyRun), }) }) /** Waits for one Apify Actor run to finish or for a caller-defined limit. */ export const waitForApifyRun = defineAction("Wait for Apify Actor run") .describe("Waits for an Apify Actor run and returns its latest state.") .account("apify") .input( z.object({ runId: z.string().min(1), /** Maximum client-side wait in seconds. Omit to wait indefinitely. */ waitSecs: z.number().int().positive().optional(), }), ) .output(APIFY_RUN_SCHEMA) .retry({ replaySafety: "safe" }) .handler(async ({ account, input }) => mapApifyRun( await getApifyClient(account.secret) .run(input.runId) .waitForFinish({ waitSecs: input.waitSecs }), ), ) /** Aborts an Apify Actor run. */ export const abortApifyRun = defineAction("Abort Apify Actor run") .describe("Stops a running Apify Actor run immediately or gracefully.") .account("apify") .input( z.object({ gracefully: z.boolean().optional(), runId: z.string().min(1), }), ) .output(APIFY_RUN_SCHEMA) .retry({ replaySafety: "unsafe" }) .handler(async ({ account, input }) => mapApifyRun( await getApifyClient(account.secret) .run(input.runId) .abort({ gracefully: input.gracefully }), ), ) /** Resurrects a finished Apify Actor run with optional resource overrides. */ export const resurrectApifyRun = defineAction("Resurrect Apify Actor run") .describe( "Starts a new run from a finished Apify Actor run and its storages.", ) .account("apify") .input(APIFY_RUN_OPTIONS_SCHEMA.extend({ runId: z.string().min(1) })) .output(APIFY_RUN_SCHEMA) .retry({ replaySafety: "unsafe" }) .handler(async ({ account, input }) => { const { runId, ...options } = input return mapApifyRun( await getApifyClient(account.secret) .run(runId) .resurrect(mapRunOptions(options)), ) }) /** Lists structured items from an Apify dataset. */ export const listApifyDatasetItems = defineAction("List Apify dataset items") .describe("Reads one page of structured items from an Apify dataset.") .account("apify") .input( z.object({ clean: z.boolean().optional(), datasetId: z.string().min(1), desc: z.boolean().optional(), fields: z.string().min(1).array().min(1).optional(), flatten: z.string().min(1).array().min(1).optional(), limit: z.number().int().positive().optional(), offset: z.number().int().nonnegative().optional(), omit: z.string().min(1).array().min(1).optional(), skipEmpty: z.boolean().optional(), skipHidden: z.boolean().optional(), unwind: z .union([z.string().min(1), z.string().min(1).array().min(1)]) .optional(), view: z.string().min(1).optional(), }), ) .output(APIFY_DATASET_ITEMS_OUTPUT_SCHEMA) .retry({ replaySafety: "safe" }) .handler(async ({ account, input }) => { const { datasetId, ...options } = input return APIFY_DATASET_ITEMS_OUTPUT_SCHEMA.parse( await getApifyClient(account.secret) .dataset>(datasetId) .listItems(options), ) }) /** Gets a JSON or text record from an Apify key-value store. */ export const getApifyKeyValueStoreRecord = defineAction( "Get Apify key-value store record", ) .describe("Gets a JSON or text value from an Apify key-value store.") .account("apify") .input( z.object({ key: z.string().min(1), storeId: z.string().min(1), }), ) .output( z.object({ contentType: z.string().optional(), key: z.string(), value: z.json(), }), ) .retry({ replaySafety: "safe" }) .handler(async ({ account, input }) => { const record = await requireResource( getApifyClient(account.secret) .keyValueStore(input.storeId) .getRecord(input.key), "key-value store record", input.key, ) return { contentType: record.contentType, key: record.key, value: z.json().parse(record.value), } }) /** Gets text from an Apify Actor run log. */ export const getApifyRunLog = defineAction("Get Apify Actor run log") .describe("Gets processed or raw text from an Apify Actor run log.") .account("apify") .input( z.object({ raw: z.boolean().optional(), runId: z.string().min(1), }), ) .output(z.object({ log: z.string() })) .retry({ replaySafety: "safe" }) .handler(async ({ account, input }) => ({ log: await requireResource( getApifyClient(account.secret) .run(input.runId) .log() .get({ raw: input.raw }), "Actor run log", input.runId, ), })) /** Starts an Apify Actor build. */ export const startApifyActorBuild = defineAction("Start Apify Actor build") .describe("Starts a build for one Apify Actor version.") .account("apify") .input( z.object({ actorId: z.string().min(1), betaPackages: z.boolean().optional(), tag: z.string().min(1).optional(), useCache: z.boolean().optional(), versionNumber: z.string().min(1), }), ) .output(APIFY_BUILD_SCHEMA) .retry({ replaySafety: "unsafe" }) .handler(async ({ account, input }) => { const { actorId, versionNumber, ...options } = input return mapApifyBuild( await getApifyClient(account.secret) .actor(actorId) .build(versionNumber, { ...options, waitForFinish: 0, }), ) }) /** Gets one Apify Actor build. */ export const getApifyActorBuild = defineAction("Get Apify Actor build") .describe("Gets the current state of one Apify Actor build.") .account("apify") .input(z.object({ buildId: z.string().min(1) })) .output(APIFY_BUILD_SCHEMA) .retry({ replaySafety: "safe" }) .handler(async ({ account, input }) => mapApifyBuild( await requireResource( getApifyClient(account.secret).build(input.buildId).get(), "Actor build", input.buildId, ), ), ) /** Waits for one Apify Actor build to finish or for a caller-defined limit. */ export const waitForApifyActorBuild = defineAction("Wait for Apify Actor build") .describe("Waits for an Apify Actor build and returns its latest state.") .account("apify") .input( z.object({ buildId: z.string().min(1), waitSecs: z.number().int().positive().optional(), }), ) .output(APIFY_BUILD_SCHEMA) .retry({ replaySafety: "safe" }) .handler(async ({ account, input }) => mapApifyBuild( await getApifyClient(account.secret) .build(input.buildId) .waitForFinish({ waitSecs: input.waitSecs }), ), ) /** Aborts an Apify Actor build. */ export const abortApifyActorBuild = defineAction("Abort Apify Actor build") .describe("Stops a running Apify Actor build.") .account("apify") .input(z.object({ buildId: z.string().min(1) })) .output(APIFY_BUILD_SCHEMA) .retry({ replaySafety: "unsafe" }) .handler(async ({ account, input }) => mapApifyBuild( await getApifyClient(account.secret).build(input.buildId).abort(), ), ) /** * Maps public run options to the official client names. * * @param options - Validated public run options. */ function mapRunOptions(options: z.output) { const { memoryMbytes, timeoutSecs, ...rest } = options return { ...rest, memory: memoryMbytes, timeout: timeoutSecs } } /** * Returns whether a memory allocation is a power of two. * * @param value - Memory in megabytes. */ function isPowerOfTwo(value: number) { return Number.isInteger(Math.log2(value)) } /** * Requires one optional provider resource. * * @param resource - Pending provider lookup. * @param name - Resource name used in the error. * @param id - Requested provider ID. */ async function requireResource( resource: Promise, name: string, id: string, ) { const value = await resource if (value === undefined) throw new Error(`Apify ${name} ${id} was not found.`) return value }