Skip to main content

ironflow_ops_gitlab/
paged_operation.rs

1//! [`GitLabPagedOp`] -- executes a paginated [`Endpoint`] across every requested page as a
2//! tracked [`Operation`].
3
4use async_trait::async_trait;
5use gitlab::AsyncGitlab;
6use gitlab::api::{AsyncQuery, Endpoint, Pageable, Paged, Pagination, paged};
7use ironflow_core::error::OperationError;
8use ironflow_core::operation::{Operation, OperationContext};
9use serde_json::Value;
10
11/// A paginated GitLab endpoint wrapped as an Ironflow [`Operation`].
12///
13/// Created via [`GitLab::paged_op`](crate::GitLab::paged_op). Unlike [`GitLabOp`](crate::GitLabOp),
14/// which issues a single request, this drives pagination itself and concatenates every page
15/// requested by `pagination` into a single JSON array.
16///
17/// # Examples
18///
19/// ```no_run
20/// use ironflow_ops_gitlab::GitLab;
21/// use gitlab::api::projects::merge_requests::MergeRequests;
22/// use gitlab::api::Pagination;
23///
24/// # async fn example() -> Result<(), ironflow_core::error::OperationError> {
25/// let gitlab = GitLab::new("glpat-xxxx", "gitlab.com").await?;
26/// let endpoint = MergeRequests::builder().project(42).build().unwrap();
27/// let op = gitlab.paged_op(endpoint, Pagination::All);
28/// # Ok(())
29/// # }
30/// ```
31///
32/// # Errors
33///
34/// [`Operation::execute`] returns [`OperationError::Http`] if any page request fails.
35pub struct GitLabPagedOp<E> {
36    client: AsyncGitlab,
37    paged: Paged<E>,
38    endpoint_path: String,
39    pagination: Pagination,
40}
41
42impl<E> GitLabPagedOp<E>
43where
44    E: Endpoint,
45{
46    pub(crate) fn new(client: AsyncGitlab, endpoint: E, pagination: Pagination) -> Self {
47        let endpoint_path = endpoint.endpoint().into_owned();
48        let paged = paged(endpoint, pagination);
49        Self {
50            client,
51            paged,
52            endpoint_path,
53            pagination,
54        }
55    }
56}
57
58#[async_trait]
59impl<E> Operation for GitLabPagedOp<E>
60where
61    E: Endpoint + Pageable + Sync + Send,
62{
63    fn kind(&self) -> &str {
64        "gitlab"
65    }
66
67    async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
68        let results: Vec<Value> =
69            self.paged
70                .query_async(&self.client)
71                .await
72                .map_err(|e| OperationError::Http {
73                    status: None,
74                    message: e.to_string(),
75                })?;
76        Ok(Value::Array(results))
77    }
78
79    fn input(&self) -> Option<Value> {
80        let pagination = match self.pagination {
81            Pagination::All => "all".to_string(),
82            Pagination::AllPerPageLimit(n) => format!("all_per_page_limit({n})"),
83            Pagination::Limit(n) => format!("limit({n})"),
84            _ => "unknown".to_string(),
85        };
86        Some(Value::Object(serde_json::Map::from_iter([
87            (
88                "endpoint".to_string(),
89                Value::String(self.endpoint_path.clone()),
90            ),
91            ("pagination".to_string(), Value::String(pagination)),
92        ])))
93    }
94}