Skip to main content

openapp_sdk_core/resources/
scripting.rs

1//! `OpenApp` Scripting resource group.
2//!
3//! Scripting is an asynchronous job API: `POST /scripting/executions` accepts a
4//! program and returns a job, which the caller polls until `status` reaches a
5//! terminal value. [`ScriptingClient::execute`] wraps that submit-then-poll
6//! cycle for callers that just want the result.
7
8use std::sync::Arc;
9use std::time::Duration;
10
11use anyhow::anyhow;
12use reqwest::Method;
13
14use super::JsonValue;
15use crate::{
16    error::SdkError,
17    transport::{RequestSpec, Transport},
18};
19
20/// How [`ScriptingClient::execute`] polls a job to completion.
21///
22/// The defaults follow the backend contract: poll every second, back off to
23/// five seconds, and give up at the 15-minute server-side runtime cap (an
24/// execution that runs longer is failed server-side, so waiting past it can
25/// never observe a result). Override via [`ScriptingClient::with_poll_schedule`]
26/// when a caller needs a different budget.
27#[derive(Debug, Clone, Copy, PartialEq, Eq)]
28pub struct PollSchedule {
29    /// Wait before the first poll.
30    pub initial_delay: Duration,
31    /// Upper bound the delay backs off to.
32    pub max_delay: Duration,
33    /// Give up with an error once this much wall time has elapsed without a
34    /// terminal status.
35    pub timeout: Duration,
36}
37
38impl PollSchedule {
39    /// Build a schedule from its three parts.
40    #[must_use]
41    pub const fn new(initial_delay: Duration, max_delay: Duration, timeout: Duration) -> Self {
42        Self {
43            initial_delay,
44            max_delay,
45            timeout,
46        }
47    }
48}
49
50impl Default for PollSchedule {
51    fn default() -> Self {
52        Self {
53            initial_delay: Duration::from_secs(1),
54            max_delay: Duration::from_secs(5),
55            timeout: Duration::from_mins(15),
56        }
57    }
58}
59
60#[derive(Debug, Clone)]
61pub struct ScriptingClient {
62    transport: Arc<Transport>,
63    poll: PollSchedule,
64}
65
66impl ScriptingClient {
67    pub(crate) fn new(transport: Arc<Transport>) -> Self {
68        Self {
69            transport,
70            poll: PollSchedule::default(),
71        }
72    }
73
74    /// Override how [`Self::execute`] polls for a terminal status.
75    #[must_use]
76    pub fn with_poll_schedule(mut self, poll: PollSchedule) -> Self {
77        self.poll = poll;
78        self
79    }
80
81    /// `POST /scripting/executions` — submit a program; returns the `202` job.
82    pub async fn create_execution(&self, body: &JsonValue) -> Result<JsonValue, SdkError> {
83        self.transport
84            .request_json::<JsonValue, JsonValue>(RequestSpec {
85                method: Method::POST,
86                path: "/scripting/executions",
87                body: Some(body),
88                ..Default::default()
89            })
90            .await
91    }
92
93    /// `GET /scripting/executions/{id}` — the only route that carries `result`.
94    pub async fn get_execution(&self, id: &str) -> Result<JsonValue, SdkError> {
95        let path = format!("/scripting/executions/{id}");
96        self.transport
97            .request_json::<(), JsonValue>(RequestSpec {
98                method: Method::GET,
99                path: &path,
100                ..Default::default()
101            })
102            .await
103    }
104
105    /// `GET /scripting/executions?limit=…` — summaries (no `result`), newest first.
106    pub async fn list_executions(&self, limit: Option<i64>) -> Result<JsonValue, SdkError> {
107        let query = [("limit", limit.map(|n| n.to_string()))];
108        self.transport
109            .request_json::<(), JsonValue>(RequestSpec {
110                method: Method::GET,
111                path: "/scripting/executions",
112                query: &query,
113                ..Default::default()
114            })
115            .await
116    }
117
118    /// `DELETE /scripting/executions/{id}` — request cooperative cancellation.
119    ///
120    /// Cancellation is not a rollback: effects the program already applied stay
121    /// applied, and the worker only notices within a few seconds.
122    pub async fn cancel_execution(&self, id: &str) -> Result<(), SdkError> {
123        let path = format!("/scripting/executions/{id}");
124        self.transport
125            .request_json::<(), ()>(RequestSpec {
126                method: Method::DELETE,
127                path: &path,
128                ..Default::default()
129            })
130            .await
131    }
132
133    /// Submit a program and wait for it to finish, returning its `result`.
134    ///
135    /// Polls [`Self::get_execution`] every `initial_delay`, doubling up to
136    /// `max_delay`, until the job is terminal or `timeout` elapses — see
137    /// [`PollSchedule`]. A script fault is **not** an HTTP error: the submission
138    /// still succeeds and the failure lands in the job, so `failed` / `canceled`
139    /// are surfaced here as [`SdkError::Other`] carrying the job's `error` text.
140    pub async fn execute(&self, body: &JsonValue) -> Result<JsonValue, SdkError> {
141        let created = self.create_execution(body).await?;
142        let id = created
143            .get("id")
144            .and_then(JsonValue::as_str)
145            .ok_or_else(|| {
146                SdkError::Deserialize(format!(
147                    "POST /scripting/executions returned no string 'id': {created}"
148                ))
149            })?
150            .to_owned();
151
152        let deadline = tokio::time::Instant::now() + self.poll.timeout;
153        let mut delay = self.poll.initial_delay;
154        loop {
155            tokio::time::sleep(delay).await;
156
157            let job = self.get_execution(&id).await?;
158            let status = job
159                .get("status")
160                .and_then(JsonValue::as_str)
161                .unwrap_or_default();
162            if is_terminal(status) {
163                return terminal_outcome(&job, status);
164            }
165
166            if tokio::time::Instant::now() >= deadline {
167                return Err(SdkError::Transport(format!(
168                    "execution {id} did not finish within {:?}",
169                    self.poll.timeout
170                )));
171            }
172            delay = (delay * 2).min(self.poll.max_delay);
173        }
174    }
175}
176
177/// `true` once no further polling can change the job's outcome.
178fn is_terminal(status: &str) -> bool {
179    matches!(status, "succeeded" | "failed" | "canceled")
180}
181
182fn terminal_outcome(job: &JsonValue, status: &str) -> Result<JsonValue, SdkError> {
183    if status == "succeeded" {
184        return Ok(job.get("result").cloned().unwrap_or(JsonValue::Null));
185    }
186
187    let error = job
188        .get("error")
189        .and_then(JsonValue::as_str)
190        .unwrap_or("no error reported");
191    Err(SdkError::Other(anyhow!(
192        "script execution {status}: {error}"
193    )))
194}