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"), ¶ms, 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"), ¶ms, 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"), ¶ms, 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}