box_open_sdk/managers/
workflows.rs1use crate::internal::path_escape;
4use crate::runtime::{self, Error};
5
6#[derive(Clone, Debug, Default)]
8pub struct WorkflowsListOptions {
9 pub trigger_type: Option<String>,
10 pub limit: Option<i64>,
11 pub marker: Option<String>,
12}
13
14pub struct WorkflowsListPaginator {
17 manager: WorkflowsManager,
18 folder_id: String,
19 options: WorkflowsListOptions,
20 buffer: std::vec::IntoIter<crate::models::schemas::Workflow>,
21 done: bool,
22}
23
24impl WorkflowsListPaginator {
25 pub async fn next(&mut self) -> Option<Result<crate::models::schemas::Workflow, Error>> {
28 loop {
29 if let Some(item) = self.buffer.next() {
30 return Some(Ok(item));
31 }
32 if self.done {
33 return None;
34 }
35 let page = match self
36 .manager
37 .list_page(self.folder_id.clone(), Some(self.options.clone()))
38 .await
39 {
40 Ok(page) => page,
41 Err(err) => {
42 self.done = true;
43 return Some(Err(err));
44 }
45 };
46 self.buffer = page.entries.unwrap_or_default().into_iter();
47 match page.next_marker.flatten() {
48 Some(cursor) if !cursor.is_empty() => self.options.marker = Some(cursor),
49 _ => self.done = true,
50 }
51 }
52 }
53}
54
55pub struct WorkflowsManager {
57 session: std::sync::Arc<runtime::Client>,
58}
59
60impl WorkflowsManager {
61 pub(crate) fn new(session: std::sync::Arc<runtime::Client>) -> Self {
62 Self { session }
63 }
64
65 async fn list_page(
66 &self,
67 folder_id: String,
68 opts: Option<WorkflowsListOptions>,
69 ) -> Result<crate::models::schemas::Workflows, Error> {
70 let mut url = self.session.base_url("api");
71 url.push_str("/workflows");
72 let mut req = self.session.new_request("GET", &url);
73 req = runtime::with_query(req, "folder_id", &folder_id);
74 let opts = opts.unwrap_or_default();
75 if let Some(value) = opts.trigger_type {
76 req = runtime::with_query(req, "trigger_type", &value);
77 }
78 if let Some(value) = opts.limit {
79 req = runtime::with_query(req, "limit", &value.to_string());
80 }
81 if let Some(value) = opts.marker {
82 req = runtime::with_query(req, "marker", &value);
83 }
84 let resp = self.session.fetch(req).await?;
85 let data = runtime::response_bytes(&resp)?;
86 Ok(serde_json::from_slice(&data)?)
87 }
88
89 pub fn list(
91 &self,
92 folder_id: String,
93 opts: Option<WorkflowsListOptions>,
94 ) -> WorkflowsListPaginator {
95 WorkflowsListPaginator {
96 manager: WorkflowsManager::new(self.session.clone()),
97 folder_id,
98 options: opts.unwrap_or_default(),
99 buffer: Vec::new().into_iter(),
100 done: false,
101 }
102 }
103
104 pub async fn start(
105 &self,
106 workflow_id: String,
107 body: crate::models::schemas::StartWorkflowRequest,
108 ) -> Result<(), Error> {
109 let mut url = self.session.base_url("api");
110 url.push_str("/workflows");
111 url.push('/');
112 let seg = path_escape(&workflow_id);
113 url.push_str(&seg);
114 url.push_str("/start");
115 let mut req = self.session.new_request("POST", &url);
116 let payload = serde_json::to_vec(&body)?;
117 req = runtime::with_json_body(req, &payload);
118 let _ = self.session.fetch(req).await?;
119 Ok(())
120 }
121}