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 ([`Task::public_config`]) set up. An Actor can have at most 10 published
54    /// tasks, and an account at most 100 across all its Actors; contact Apify support to raise
55    /// either limit. Publishing an already-published task does nothing. The returned
56    /// [`Task::is_public`] reflects the new publication state.
57    pub async fn publish(&self) -> ApifyClientResult<Task> {
58        self.update(&serde_json::json!({ "isPublic": true })).await
59    }
60
61    /// Unpublishes the task from its public landing page, by setting `isPublic: false` through
62    /// [`TaskClient::update`].
63    ///
64    /// The public display configuration ([`Task::public_config`]) is preserved, so the task can
65    /// be published again later without re-entering it. Unpublishing a task that is not
66    /// published does nothing.
67    pub async fn unpublish(&self) -> ApifyClientResult<Task> {
68        self.update(&serde_json::json!({ "isPublic": false })).await
69    }
70
71    /// Starts the task and returns immediately with the created run.
72    ///
73    /// `input` overrides the task's saved input (or `None` to use the saved input).
74    pub async fn start<T: Serialize>(
75        &self,
76        input: Option<&T>,
77        options: ActorStartOptions,
78    ) -> ApifyClientResult<ActorRun> {
79        let mut params = QueryParams::new();
80        options.apply(&mut params);
81        let body = match input {
82            Some(value) => Some(serde_json::to_vec(value)?),
83            None => None,
84        };
85        post_with_body(&self.ctx, Some("runs"), &params, body, "application/json").await
86    }
87
88    /// Starts the task and waits (client-side polling) for it to finish.
89    ///
90    /// `wait_secs` controls the wait budget:
91    /// - `None` polls indefinitely until the run reaches a terminal state.
92    /// - `Some(n)` bounds the wait to roughly `n` seconds; if the run has not finished by
93    ///   then, the **last fetched (still non-terminal) run is returned** rather than an
94    ///   error. Check `status` / `is_terminal()` on the result when using `Some`.
95    pub async fn call<T: Serialize>(
96        &self,
97        input: Option<&T>,
98        options: ActorStartOptions,
99        wait_secs: Option<i64>,
100    ) -> ApifyClientResult<ActorRun> {
101        let run = self.start(input, options).await?;
102        self.root.run(run.id).wait_for_finish(wait_secs).await
103    }
104
105    /// Fetches the task's saved input, or `None` if not set.
106    pub async fn get_input(&self) -> ApifyClientResult<Option<Value>> {
107        let response =
108            crate::clients::base::get_raw(&self.ctx, Some("input"), &QueryParams::new()).await?;
109        match response {
110            Some(r) => Ok(Some(serde_json::from_slice(&r.body)?)),
111            None => Ok(None),
112        }
113    }
114
115    /// Updates the task's saved input.
116    pub async fn update_input<T: Serialize>(&self, input: &T) -> ApifyClientResult<Value> {
117        let body = serde_json::to_vec(input)?;
118        let url = self.ctx.url(Some("input"));
119        let mut headers = std::collections::HashMap::new();
120        headers.insert("Content-Type".to_string(), "application/json".to_string());
121        let response = self
122            .ctx
123            .http
124            .call(crate::http_client::HttpRequest {
125                method: crate::http_client::HttpMethod::Put,
126                url,
127                headers,
128                body: Some(body),
129                timeout: crate::clients::base::DEFAULT_REQUEST_TIMEOUT,
130            })
131            .await?;
132        Ok(serde_json::from_slice(&response.body)?)
133    }
134
135    /// Returns a client for the last run of this task, optionally filtered by run status.
136    ///
137    /// `status` filters by run status (e.g. `"SUCCEEDED"`, `"FAILED"`, `"RUNNING"`); pass `None`
138    /// to leave it unfiltered. This maps to the `status` query parameter on
139    /// `GET /v2/actor-tasks/{actorTaskId}/runs/last` and mirrors the reference client's
140    /// `lastRun({ status })`. To also filter by `origin`, use [`TaskClient::last_run_with_options`].
141    pub fn last_run(&self, status: Option<&str>) -> RunClient {
142        self.last_run_with_options(LastRunOptions {
143            status: status.map(str::to_owned),
144            origin: None,
145        })
146    }
147
148    /// Returns a client for the last run of this task, applying the given [`LastRunOptions`]
149    /// (e.g. [`LastRunOptions::status`] and/or [`LastRunOptions::origin`]).
150    ///
151    /// `status` filters by run status (e.g. `"SUCCEEDED"`, `"FAILED"`, `"RUNNING"`); `origin` filters
152    /// by how the run was started, with accepted values being the platform's run origins (e.g.
153    /// `"DEVELOPMENT"`, `"WEB"`, `"API"`, `"SCHEDULER"`). Both are documented optional query
154    /// parameters on `GET /v2/actor-tasks/{actorTaskId}/runs/last` and match the reference client's
155    /// `lastRun({ status, origin })`; leave a field as `None` to omit it.
156    pub fn last_run_with_options(&self, options: LastRunOptions) -> RunClient {
157        let mut client = RunClient::new(
158            self.root.clone(),
159            self.ctx.http.clone(),
160            &self.ctx.url(None),
161            "runs",
162            "last",
163        );
164        if let Some(status) = options.status.as_deref() {
165            client.set_base_param("status", status);
166        }
167        if let Some(origin) = options.origin.as_deref() {
168            client.set_base_param("origin", origin);
169        }
170        client
171    }
172
173    /// Returns a client for this task's run collection.
174    pub fn runs(&self) -> RunCollectionClient {
175        RunCollectionClient::new(self.ctx.http.clone(), &self.ctx.url(None), "runs")
176    }
177
178    /// Returns a client for this task's webhook collection.
179    pub fn webhooks(&self) -> WebhookCollectionClient {
180        WebhookCollectionClient::with_base(self.ctx.http.clone(), &self.ctx.url(None))
181    }
182}