Skip to main content

rectilinear_core/linear/
cycles.rs

1use std::future::ready;
2
3use anyhow::Result;
4use serde::Deserialize;
5use uuid::Uuid;
6
7use crate::db::{self, Database};
8
9use super::pagination::{paginate, ConnectionPage, LinearOperation, PageInfo};
10use super::LinearClient;
11
12#[derive(Debug, Deserialize)]
13struct CyclesData {
14    cycles: CycleConnection,
15}
16
17#[derive(Debug, Deserialize)]
18struct CycleConnection {
19    nodes: Vec<LinearCycle>,
20    #[serde(rename = "pageInfo")]
21    page_info: PageInfo,
22}
23
24#[derive(Debug, Deserialize)]
25struct LinearCycle {
26    id: String,
27    number: i32,
28    name: Option<String>,
29    #[serde(rename = "startsAt")]
30    starts_at: Option<String>,
31    #[serde(rename = "endsAt")]
32    ends_at: Option<String>,
33    #[serde(rename = "completedAt")]
34    completed_at: Option<String>,
35    #[serde(rename = "archivedAt")]
36    archived_at: Option<String>,
37    #[serde(rename = "createdAt")]
38    created_at: String,
39    #[serde(rename = "updatedAt")]
40    updated_at: String,
41    team: LinearCycleTeam,
42}
43
44#[derive(Debug, Deserialize)]
45struct LinearCycleTeam {
46    id: String,
47    key: String,
48}
49
50impl LinearClient {
51    pub async fn sync_cycles(
52        &self,
53        db: &Database,
54        team_key: &str,
55        workspace_id: &str,
56        include_archived: bool,
57    ) -> Result<usize> {
58        let sync_token = Uuid::new_v4().to_string();
59        db.mark_sync_family_running(
60            workspace_id,
61            team_key,
62            "cycles",
63            None,
64            Some(self.sync_query_config.page_size(LinearOperation::Cycles)),
65            &sync_token,
66        )?;
67        let query = r#"
68            query($teamKey: String!, $first: Int!, $after: String, $includeArchived: Boolean!) {
69                cycles(
70                    first: $first,
71                    after: $after,
72                    filter: { team: { key: { eq: $teamKey } } },
73                    includeArchived: $includeArchived,
74                    orderBy: updatedAt
75                ) {
76                    nodes {
77                        id number name startsAt endsAt completedAt archivedAt createdAt updatedAt
78                        team { id key }
79                    }
80                    pageInfo { hasNextPage endCursor }
81                }
82            }
83        "#;
84        let result = paginate(
85            &self.sync_query_config,
86            LinearOperation::Cycles,
87            Some(team_key.to_string()),
88            |request| async move {
89                let data: CyclesData = self
90                    .query_operation(
91                        LinearOperation::Cycles.name(),
92                        request.cursor.as_deref(),
93                        query,
94                        serde_json::json!({
95                            "teamKey": team_key,
96                            "first": request.page_size,
97                            "after": request.cursor,
98                            "includeArchived": include_archived,
99                        }),
100                    )
101                    .await?;
102                Ok(ConnectionPage {
103                    nodes: data.cycles.nodes,
104                    page_info: data.cycles.page_info,
105                })
106            },
107            |nodes, _| {
108                let result = nodes.into_iter().try_for_each(|cycle| {
109                    db.upsert_cycle(
110                        &db::Cycle {
111                            id: cycle.id,
112                            workspace_id: workspace_id.to_string(),
113                            team_id: cycle.team.id,
114                            team_key: cycle.team.key,
115                            number: cycle.number,
116                            name: cycle.name,
117                            starts_at: cycle.starts_at,
118                            ends_at: cycle.ends_at,
119                            completed_at: cycle.completed_at,
120                            archived_at: cycle.archived_at,
121                            created_at: cycle.created_at,
122                            updated_at: cycle.updated_at,
123                        },
124                        &sync_token,
125                    )
126                });
127                ready(result)
128            },
129            |cycle| cycle.id.clone(),
130            |event| self.observe_sync_event(event),
131        )
132        .await;
133
134        match result {
135            Ok(stats) => {
136                db.reconcile_cycles(workspace_id, team_key, &sync_token)?;
137                db.mark_sync_family_complete(
138                    workspace_id,
139                    team_key,
140                    "cycles",
141                    Some(self.sync_query_config.page_size(LinearOperation::Cycles)),
142                    &sync_token,
143                )?;
144                Ok(stats.nodes)
145            }
146            Err(error) => {
147                let message = self.redacted_error_message(&error);
148                db.mark_sync_family_failed(
149                    workspace_id,
150                    team_key,
151                    "cycles",
152                    &sync_token,
153                    &message,
154                )?;
155                Err(error)
156            }
157        }
158    }
159}