openapp_sdk_core/resources/
scripting.rs1use 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
28pub struct PollSchedule {
29 pub initial_delay: Duration,
31 pub max_delay: Duration,
33 pub timeout: Duration,
36}
37
38impl PollSchedule {
39 #[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 #[must_use]
76 pub fn with_poll_schedule(mut self, poll: PollSchedule) -> Self {
77 self.poll = poll;
78 self
79 }
80
81 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 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 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 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 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
177fn 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}