rectilinear_core/linear/
cycles.rs1use 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}