Skip to main content

apify_client/clients/
task.rs

1//! Client for a single Actor task (`/v2/actor-tasks/{actorTaskId}`).
2
3use serde::Serialize;
4use serde_json::Value;
5
6use crate::client::ApifyClient;
7use crate::clients::actor::ActorStartOptions;
8use crate::clients::base::{
9    delete_resource, get_resource, post_with_body, update_resource, ResourceContext,
10};
11use crate::clients::run::{LastRunOptions, RunClient};
12use crate::clients::run_collection::RunCollectionClient;
13use crate::clients::webhook_collection::WebhookCollectionClient;
14use crate::common::QueryParams;
15use crate::error::ApifyClientResult;
16use crate::http_client::HttpClient;
17use crate::models::{ActorRun, Task};
18
19/// Client for a specific Actor task.
20#[derive(Debug, Clone)]
21pub struct TaskClient {
22    root: ApifyClient,
23    ctx: ResourceContext,
24}
25
26impl TaskClient {
27    pub(crate) fn new(root: ApifyClient, http: HttpClient, base_url: &str, id: &str) -> Self {
28        Self {
29            root,
30            ctx: ResourceContext::single(http, base_url, "actor-tasks", id),
31        }
32    }
33
34    /// Fetches the task object, or `None` if it does not exist.
35    pub async fn get(&self) -> ApifyClientResult<Option<Task>> {
36        get_resource(&self.ctx, None, &QueryParams::new()).await
37    }
38
39    /// Updates the task with the given fields.
40    pub async fn update<T: Serialize>(&self, new_fields: &T) -> ApifyClientResult<Task> {
41        update_resource(&self.ctx, None, new_fields).await
42    }
43
44    /// Deletes the task.
45    pub async fn delete(&self) -> ApifyClientResult<()> {
46        delete_resource(&self.ctx, None).await
47    }
48
49    /// Publishes the task on its public landing page, by setting `isPublic: true` through
50    /// [`TaskClient::update`].
51    ///
52    /// The task's Actor must be public and the task must already have its public display
53    /// configuration (`publicConfig`) set up. Publishing an already-published task does nothing.
54    pub async fn publish(&self) -> ApifyClientResult<Task> {
55        self.update(&serde_json::json!({ "isPublic": true })).await
56    }
57
58    /// Unpublishes the task from its public landing page, by setting `isPublic: false` through
59    /// [`TaskClient::update`].
60    ///
61    /// The public display configuration (`publicConfig`) is preserved, so the task can be
62    /// published again later without re-entering it. Unpublishing a task that is not published
63    /// does nothing.
64    pub async fn unpublish(&self) -> ApifyClientResult<Task> {
65        self.update(&serde_json::json!({ "isPublic": false })).await
66    }
67
68    /// Starts the task and returns immediately with the created run.
69    ///
70    /// `input` overrides the task's saved input (or `None` to use the saved input).
71    pub async fn start<T: Serialize>(
72        &self,
73        input: Option<&T>,
74        options: ActorStartOptions,
75    ) -> ApifyClientResult<ActorRun> {
76        let mut params = QueryParams::new();
77        options.apply(&mut params);
78        let body = match input {
79            Some(value) => Some(serde_json::to_vec(value)?),
80            None => None,
81        };
82        post_with_body(&self.ctx, Some("runs"), &params, body, "application/json").await
83    }
84
85    /// Starts the task and waits (client-side polling) for it to finish.
86    ///
87    /// `wait_secs` controls the wait budget:
88    /// - `None` polls indefinitely until the run reaches a terminal state.
89    /// - `Some(n)` bounds the wait to roughly `n` seconds; if the run has not finished by
90    ///   then, the **last fetched (still non-terminal) run is returned** rather than an
91    ///   error. Check `status` / `is_terminal()` on the result when using `Some`.
92    pub async fn call<T: Serialize>(
93        &self,
94        input: Option<&T>,
95        options: ActorStartOptions,
96        wait_secs: Option<i64>,
97    ) -> ApifyClientResult<ActorRun> {
98        let run = self.start(input, options).await?;
99        self.root.run(run.id).wait_for_finish(wait_secs).await
100    }
101
102    /// Fetches the task's saved input, or `None` if not set.
103    pub async fn get_input(&self) -> ApifyClientResult<Option<Value>> {
104        let response =
105            crate::clients::base::get_raw(&self.ctx, Some("input"), &QueryParams::new()).await?;
106        match response {
107            Some(r) => Ok(Some(serde_json::from_slice(&r.body)?)),
108            None => Ok(None),
109        }
110    }
111
112    /// Updates the task's saved input.
113    pub async fn update_input<T: Serialize>(&self, input: &T) -> ApifyClientResult<Value> {
114        let body = serde_json::to_vec(input)?;
115        let url = self.ctx.url(Some("input"));
116        let mut headers = std::collections::HashMap::new();
117        headers.insert("Content-Type".to_string(), "application/json".to_string());
118        let response = self
119            .ctx
120            .http
121            .call(crate::http_client::HttpRequest {
122                method: crate::http_client::HttpMethod::Put,
123                url,
124                headers,
125                body: Some(body),
126                timeout: crate::clients::base::DEFAULT_REQUEST_TIMEOUT,
127            })
128            .await?;
129        Ok(serde_json::from_slice(&response.body)?)
130    }
131
132    /// Returns a client for the last run of this task, optionally filtered by run status.
133    ///
134    /// `status` filters by run status (e.g. `"SUCCEEDED"`, `"FAILED"`, `"RUNNING"`); pass `None`
135    /// to leave it unfiltered. This maps to the `status` query parameter on
136    /// `GET /v2/actor-tasks/{actorTaskId}/runs/last` and mirrors the reference client's
137    /// `lastRun({ status })`. To also filter by `origin`, use [`TaskClient::last_run_with_options`].
138    pub fn last_run(&self, status: Option<&str>) -> RunClient {
139        self.last_run_with_options(LastRunOptions {
140            status: status.map(str::to_owned),
141            origin: None,
142        })
143    }
144
145    /// Returns a client for the last run of this task, applying the given [`LastRunOptions`]
146    /// (e.g. [`LastRunOptions::status`] and/or [`LastRunOptions::origin`]).
147    ///
148    /// `status` filters by run status (e.g. `"SUCCEEDED"`, `"FAILED"`, `"RUNNING"`); `origin` filters
149    /// by how the run was started, with accepted values being the platform's run origins (e.g.
150    /// `"DEVELOPMENT"`, `"WEB"`, `"API"`, `"SCHEDULER"`). Both are documented optional query
151    /// parameters on `GET /v2/actor-tasks/{actorTaskId}/runs/last` and match the reference client's
152    /// `lastRun({ status, origin })`; leave a field as `None` to omit it.
153    pub fn last_run_with_options(&self, options: LastRunOptions) -> RunClient {
154        let mut client = RunClient::new(
155            self.root.clone(),
156            self.ctx.http.clone(),
157            &self.ctx.url(None),
158            "runs",
159            "last",
160        );
161        if let Some(status) = options.status.as_deref() {
162            client.set_base_param("status", status);
163        }
164        if let Some(origin) = options.origin.as_deref() {
165            client.set_base_param("origin", origin);
166        }
167        client
168    }
169
170    /// Returns a client for this task's run collection.
171    pub fn runs(&self) -> RunCollectionClient {
172        RunCollectionClient::new(self.ctx.http.clone(), &self.ctx.url(None), "runs")
173    }
174
175    /// Returns a client for this task's webhook collection.
176    pub fn webhooks(&self) -> WebhookCollectionClient {
177        WebhookCollectionClient::with_base(self.ctx.http.clone(), &self.ctx.url(None))
178    }
179}