Skip to main content

apify_client/clients/
run.rs

1//! Client for a single Actor run (`/v2/actor-runs/{runId}` and nested routes).
2
3use serde::Serialize;
4
5use crate::clients::base::{
6    delete_resource, get_resource, post_action, post_with_body, update_resource, wait_for_finish,
7    ResourceContext,
8};
9use crate::clients::dataset::DatasetClient;
10use crate::clients::key_value_store::KeyValueStoreClient;
11use crate::clients::log::LogClient;
12use crate::clients::request_queue::RequestQueueClient;
13use crate::common::{to_safe_id, QueryParams};
14use crate::error::ApifyClientResult;
15use crate::http_client::HttpClient;
16use crate::models::ActorRun;
17
18/// Header the API uses to deduplicate charge requests (matching the reference client).
19const CHARGE_IDEMPOTENCY_HEADER: &str = "idempotency-key";
20
21/// Options for fetching an Actor's or task's last run
22/// ([`crate::clients::actor::ActorClient::last_run_with_options`] /
23/// [`crate::clients::task::TaskClient::last_run_with_options`]).
24///
25/// Both are sent as query parameters to the last-run endpoints
26/// (`GET /v2/actors/{actorId}/runs/last`, `GET /v2/actor-tasks/{actorTaskId}/runs/last`).
27/// `status` and `origin` are both documented optional filters on those endpoints and match the
28/// reference client's `lastRun({ status, origin })`.
29#[derive(Debug, Default, Clone)]
30pub struct LastRunOptions {
31    /// Filter by run status (e.g. `"SUCCEEDED"`, `"FAILED"`, `"RUNNING"`). `None` leaves it
32    /// unfiltered.
33    pub status: Option<String>,
34    /// Filter by how the run was started; accepted values are the platform's run origins
35    /// (e.g. `"DEVELOPMENT"`, `"WEB"`, `"API"`, `"SCHEDULER"`). `None` leaves it unfiltered.
36    pub origin: Option<String>,
37}
38
39/// Options for resurrecting a finished run.
40#[derive(Debug, Default, Clone)]
41pub struct RunResurrectOptions {
42    /// Build tag/number to use; defaults to the original run's build.
43    pub build: Option<String>,
44    /// Memory in megabytes.
45    pub memory_mbytes: Option<i64>,
46    /// Timeout in seconds.
47    pub timeout_secs: Option<i64>,
48    /// Maximum number of dataset items to charge (pay-per-result Actors).
49    pub max_items: Option<i64>,
50    /// Maximum total charge in USD (pay-per-event Actors).
51    pub max_total_charge_usd: Option<f64>,
52    /// If `true`, restart the run automatically when it fails.
53    pub restart_on_error: Option<bool>,
54}
55
56/// Options for transforming a run into another Actor's run (metamorph).
57#[derive(Debug, Default, Clone)]
58pub struct RunMetamorphOptions {
59    /// Build tag/number of the target Actor to use (defaults to the target's default build).
60    pub build: Option<String>,
61    /// Content type of the input body. Defaults to `application/json` when unset.
62    pub content_type: Option<String>,
63}
64
65/// Options for charging a pay-per-event run via [`RunClient::charge`].
66#[derive(Debug, Default, Clone)]
67pub struct RunChargeOptions {
68    /// Name of the event to charge for. Required.
69    pub event_name: String,
70    /// Number of times to charge the event (defaults to `1`).
71    pub count: Option<i64>,
72    /// Idempotency key deduplicating the charge across retries. The API requires this header
73    /// on every charge request and forgets each key 3 minutes after the charge, so a retry sent
74    /// past that window with the same key creates a new charge rather than being deduplicated.
75    /// If `None`, one is auto-generated as `{runId}-{eventName}-{timestampMillis}-{random}`,
76    /// matching the reference client, so a transport-retried charge is applied at most once.
77    pub idempotency_key: Option<String>,
78}
79
80/// Client for a specific Actor run.
81///
82/// In addition to CRUD-style operations, this client exposes the run's lifecycle
83/// actions (abort, metamorph, reboot, resurrect, charge) and provides access to the
84/// run's default dataset, key-value store, request queue and log.
85#[derive(Debug, Clone)]
86pub struct RunClient {
87    ctx: ResourceContext,
88    /// The run ID, retained so `charge` can build a per-run idempotency key.
89    id: String,
90}
91
92impl RunClient {
93    pub(crate) fn new(
94        _root: crate::client::ApifyClient,
95        http: HttpClient,
96        base_url: &str,
97        resource_path: &str,
98        id: &str,
99    ) -> Self {
100        Self {
101            ctx: ResourceContext::single(http, base_url, resource_path, id),
102            id: id.to_string(),
103        }
104    }
105
106    /// Adds a base query parameter to this run client, inherited by its requests (used by
107    /// `actor.last_run` / `task.last_run` to thread the `status` and `origin` filters).
108    pub(crate) fn set_base_param(&mut self, key: &str, value: &str) {
109        self.ctx
110            .base_params
111            .push_raw(key.to_string(), value.to_string());
112    }
113
114    /// Fetches the run object, or `None` if it does not exist.
115    pub async fn get(&self) -> ApifyClientResult<Option<ActorRun>> {
116        get_resource(&self.ctx, None, &QueryParams::new()).await
117    }
118
119    /// Updates the run (e.g. its status message) and returns the updated object.
120    pub async fn update<T: Serialize>(&self, new_fields: &T) -> ApifyClientResult<ActorRun> {
121        update_resource(&self.ctx, None, new_fields).await
122    }
123
124    /// Deletes the run.
125    pub async fn delete(&self) -> ApifyClientResult<()> {
126        delete_resource(&self.ctx, None).await
127    }
128
129    /// Aborts the run. `gracefully` is optional, matching the reference client's optional
130    /// `gracefully` option and the Go sibling's `Option<bool>`: `Some(true)` lets the run
131    /// perform cleanup first, `Some(false)` aborts immediately, and `None` omits the parameter
132    /// entirely so the server applies its default (immediate abort).
133    pub async fn abort(&self, gracefully: Option<bool>) -> ApifyClientResult<ActorRun> {
134        let mut params = QueryParams::new();
135        params.add_bool("gracefully", gracefully);
136        post_action(&self.ctx, Some("abort"), &params, None, None).await
137    }
138
139    /// Transforms the run into a run of another Actor (metamorph).
140    ///
141    /// `options.content_type` sets the content type of the input body (defaulting to
142    /// `application/json`), matching the reference client's `metamorph(..., { contentType })`.
143    /// To send a non-JSON input as raw bytes instead, use [`metamorph_raw`](Self::metamorph_raw).
144    pub async fn metamorph<T: Serialize>(
145        &self,
146        target_actor_id: &str,
147        input: Option<&T>,
148        options: RunMetamorphOptions,
149    ) -> ApifyClientResult<ActorRun> {
150        let body = match input {
151            Some(value) => Some(serde_json::to_vec(value)?),
152            None => None,
153        };
154        self.metamorph_with_body(target_actor_id, body, "application/json", options)
155            .await
156    }
157
158    /// Transforms the run into a run of another Actor (metamorph) with a raw request body
159    /// instead of a JSON-serializable value.
160    ///
161    /// The bytes are sent exactly as given, with no JSON serialization — pair this with
162    /// `options.content_type` for a non-JSON input. For an object or array input, use
163    /// [`metamorph`](Self::metamorph) instead.
164    pub async fn metamorph_raw(
165        &self,
166        target_actor_id: &str,
167        input: &[u8],
168        options: RunMetamorphOptions,
169    ) -> ApifyClientResult<ActorRun> {
170        self.metamorph_with_body(
171            target_actor_id,
172            Some(input.to_vec()),
173            "application/octet-stream",
174            options,
175        )
176        .await
177    }
178
179    /// Shared implementation of [`metamorph`](Self::metamorph) and
180    /// [`metamorph_raw`](Self::metamorph_raw).
181    async fn metamorph_with_body(
182        &self,
183        target_actor_id: &str,
184        body: Option<Vec<u8>>,
185        default_content_type: &str,
186        options: RunMetamorphOptions,
187    ) -> ApifyClientResult<ActorRun> {
188        let mut params = QueryParams::new();
189        params
190            .add_str("targetActorId", Some(to_safe_id(target_actor_id)))
191            .add_str("build", options.build);
192        let content_type = options
193            .content_type
194            .as_deref()
195            .unwrap_or(default_content_type)
196            .to_string();
197        post_with_body(&self.ctx, Some("metamorph"), &params, body, &content_type).await
198    }
199
200    /// Reboots the run (restarts its container, preserving the run ID and storages).
201    pub async fn reboot(&self) -> ApifyClientResult<ActorRun> {
202        post_action(&self.ctx, Some("reboot"), &QueryParams::new(), None, None).await
203    }
204
205    /// Resurrects a finished run, starting it again with (optionally overridden) settings.
206    pub async fn resurrect(&self, options: RunResurrectOptions) -> ApifyClientResult<ActorRun> {
207        let mut params = QueryParams::new();
208        params
209            .add_str("build", options.build)
210            .add_int("memory", options.memory_mbytes)
211            .add_int("timeout", options.timeout_secs)
212            .add_int("maxItems", options.max_items)
213            .add_float("maxTotalChargeUsd", options.max_total_charge_usd)
214            .add_bool("restartOnError", options.restart_on_error);
215        post_action(&self.ctx, Some("resurrect"), &params, None, None).await
216    }
217
218    /// Charges the run for a pay-per-event run, recording occurrences of a named event.
219    ///
220    /// An idempotency key is always sent (auto-generated when `options.idempotency_key` is
221    /// `None`), so a charge that is retried by the transport is applied at most once — matching
222    /// the reference client and preventing double-charging.
223    ///
224    /// The charge endpoint returns an empty body on success, so this issues the request
225    /// directly and treats any 2xx response as success (errors still surface normally).
226    pub async fn charge(&self, options: RunChargeOptions) -> ApifyClientResult<()> {
227        let count = options.count.unwrap_or(1);
228        let idempotency_key = options
229            .idempotency_key
230            .unwrap_or_else(|| self.generate_idempotency_key(&options.event_name));
231        let body = serde_json::json!({ "eventName": options.event_name, "count": count });
232        let body_bytes = serde_json::to_vec(&body)?;
233        let url = self.ctx.url(Some("charge"));
234        let mut headers = std::collections::HashMap::new();
235        headers.insert("Content-Type".to_string(), "application/json".to_string());
236        headers.insert(CHARGE_IDEMPOTENCY_HEADER.to_string(), idempotency_key);
237        // A successful `HttpClient::call` already guarantees a 2xx status; the (empty) body
238        // is intentionally ignored rather than parsed as a `data` envelope.
239        self.ctx
240            .http
241            .call(crate::http_client::HttpRequest {
242                method: crate::http_client::HttpMethod::Post,
243                url,
244                headers,
245                body: Some(body_bytes),
246                timeout: crate::clients::base::DEFAULT_REQUEST_TIMEOUT,
247            })
248            .await?;
249        Ok(())
250    }
251
252    /// Builds a per-charge idempotency key of the form
253    /// `{runId}-{eventName}-{timestampMillis}-{random}`, matching the reference client. The
254    /// suffix only needs to be unique enough to avoid collisions within the same millisecond;
255    /// it is derived from the sub-millisecond part of the current time (no crypto needed).
256    fn generate_idempotency_key(&self, event_name: &str) -> String {
257        let now = std::time::SystemTime::now()
258            .duration_since(std::time::UNIX_EPOCH)
259            .unwrap_or_default();
260        let millis = now.as_millis();
261        // Sub-millisecond nanos give a cheap, non-crypto "random" suffix (cf. JS Math.random()).
262        let random_suffix = now.subsec_nanos() % 1_000_000;
263        format!("{}-{event_name}-{millis}-{random_suffix}", self.id)
264    }
265
266    /// Waits (by client-side polling) for the run to reach a terminal state.
267    ///
268    /// `wait_secs` controls the wait budget:
269    /// - `None` polls indefinitely until the run reaches a terminal state.
270    /// - `Some(n)` bounds the wait to roughly `n` seconds; if the run has not finished by
271    ///   then, the **last fetched (still non-terminal) run is returned** rather than an
272    ///   error. Check `status` / `is_terminal()` on the result when using `Some`.
273    pub async fn wait_for_finish(&self, wait_secs: Option<i64>) -> ApifyClientResult<ActorRun> {
274        wait_for_finish(&self.ctx, wait_secs, |r: &ActorRun| r.is_terminal()).await
275    }
276
277    /// Returns a client for the run's default dataset.
278    pub fn dataset(&self) -> DatasetClient {
279        DatasetClient::nested(self.ctx.http.clone(), &self.ctx.url(None), "dataset")
280    }
281
282    /// Returns a client for the run's default key-value store.
283    pub fn key_value_store(&self) -> KeyValueStoreClient {
284        KeyValueStoreClient::nested(
285            self.ctx.http.clone(),
286            &self.ctx.url(None),
287            "key-value-store",
288        )
289    }
290
291    /// Returns a client for the run's default request queue.
292    pub fn request_queue(&self) -> RequestQueueClient {
293        RequestQueueClient::nested(self.ctx.http.clone(), &self.ctx.url(None), "request-queue")
294    }
295
296    /// Returns a client for the run's log.
297    pub fn log(&self) -> LogClient {
298        LogClient::nested(self.ctx.http.clone(), &self.ctx.url(None), "log")
299    }
300
301    /// Opens a live stream of the run's log for redirection.
302    ///
303    /// Convenience equivalent to `run.log().stream()` (mirrors the reference client's
304    /// `getStreamedLog`): yields log chunks as they arrive, so callers can forward them to
305    /// their own logger/stdout while the run is in progress.
306    pub async fn get_streamed_log(
307        &self,
308    ) -> ApifyClientResult<impl futures_util::Stream<Item = ApifyClientResult<Vec<u8>>>> {
309        self.log().stream().await
310    }
311
312    /// Opens a live stream of the run's log for redirection, applying the given
313    /// [`LogOptions`] (e.g. [`LogOptions::raw`] to stream the unprocessed log, which is the
314    /// form the JS reference's log redirection consumes internally).
315    ///
316    /// This is a Rust-specific convenience that simply forwards `LogOptions` to
317    /// [`LogClient::stream_with_options`]; it is not a 1:1 mirror of the JS `getStreamedLog`
318    /// method's signature (which takes redirect options and returns a `StreamedLog` object).
319    pub async fn get_streamed_log_with_options(
320        &self,
321        options: crate::clients::log::LogOptions,
322    ) -> ApifyClientResult<impl futures_util::Stream<Item = ApifyClientResult<Vec<u8>>>> {
323        self.log().stream_with_options(options).await
324    }
325}