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