1use chrono::{DateTime, Utc};
20use cliban_core::time;
21use rusqlite::{params, Connection, Row};
22
23use cliban_core::{Error, Result};
24
25pub use cliban_core::migrations::REMOTE_LINKS_DDL as DDL;
29
30const COLS: &str = "id, provider, entity, local_id, remote_id, remote_key, \
31 remote_updated_at, base_hash, last_synced_at";
32
33#[derive(Debug, Clone)]
35pub struct RemoteLink {
36 pub id: i64,
37 pub provider: String,
38 pub entity: String,
39 pub local_id: i64,
40 pub remote_id: String,
44 pub remote_key: String,
46 pub remote_updated_at: Option<DateTime<Utc>>,
49 pub base_hash: Option<String>,
52 pub last_synced_at: DateTime<Utc>,
53}
54
55#[derive(Debug, Clone)]
57pub struct NewLink {
58 pub provider: String,
59 pub entity: String,
60 pub local_id: i64,
61 pub remote_id: String,
62 pub remote_key: String,
63 pub remote_updated_at: Option<DateTime<Utc>>,
64 pub base_hash: Option<String>,
65}
66
67pub fn ensure_table(conn: &Connection) -> Result<()> {
70 for stmt in DDL {
71 conn.execute_batch(stmt)?;
72 }
73 Ok(())
74}
75
76pub fn by_local(
78 conn: &Connection,
79 provider: &str,
80 entity: &str,
81 local_id: i64,
82) -> Result<Option<RemoteLink>> {
83 let sql = format!(
84 "SELECT {COLS} FROM remote_links \
85 WHERE provider = ?1 AND entity = ?2 AND local_id = ?3"
86 );
87 let mut stmt = conn.prepare(&sql)?;
88 let mut rows = stmt.query(params![provider, entity, local_id])?;
89 match rows.next()? {
90 Some(row) => Ok(Some(read(row)?)),
91 None => Ok(None),
92 }
93}
94
95pub fn by_remote(
98 conn: &Connection,
99 provider: &str,
100 entity: &str,
101 remote_id: &str,
102) -> Result<Option<RemoteLink>> {
103 let sql = format!(
104 "SELECT {COLS} FROM remote_links \
105 WHERE provider = ?1 AND entity = ?2 AND remote_id = ?3"
106 );
107 let mut stmt = conn.prepare(&sql)?;
108 let mut rows = stmt.query(params![provider, entity, remote_id])?;
109 match rows.next()? {
110 Some(row) => Ok(Some(read(row)?)),
111 None => Ok(None),
112 }
113}
114
115pub fn upsert(conn: &Connection, new: NewLink) -> Result<RemoteLink> {
122 let now = time::format_usec(time::now_usec());
123 let remote_updated = new.remote_updated_at.map(time::format_usec);
124
125 conn.execute(
126 "INSERT INTO remote_links \
127 (provider, entity, local_id, remote_id, remote_key, \
128 remote_updated_at, base_hash, last_synced_at, inserted_at, updated_at) \
129 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?8, ?8) \
130 ON CONFLICT (provider, entity, local_id) DO UPDATE SET \
131 remote_id = excluded.remote_id, \
132 remote_key = excluded.remote_key, \
133 remote_updated_at = excluded.remote_updated_at, \
134 base_hash = excluded.base_hash, \
135 last_synced_at = excluded.last_synced_at, \
136 updated_at = excluded.updated_at",
137 params![
138 new.provider,
139 new.entity,
140 new.local_id,
141 new.remote_id,
142 new.remote_key,
143 remote_updated,
144 new.base_hash,
145 now,
146 ],
147 )?;
148
149 by_local(conn, &new.provider, &new.entity, new.local_id)?.ok_or(Error::NotFound)
150}
151
152pub fn list(conn: &Connection, provider: &str) -> Result<Vec<RemoteLink>> {
154 let sql = format!("SELECT {COLS} FROM remote_links WHERE provider = ?1 ORDER BY id");
155 let mut stmt = conn.prepare(&sql)?;
156 let mut out = Vec::new();
157 let mut rows = stmt.query(params![provider])?;
158 while let Some(row) = rows.next()? {
159 out.push(read(row)?);
160 }
161 Ok(out)
162}
163
164fn read(row: &Row) -> Result<RemoteLink> {
165 let remote_updated_at: Option<String> = row.get(6)?;
166 let last_synced_at: String = row.get(8)?;
167 Ok(RemoteLink {
168 id: row.get(0)?,
169 provider: row.get(1)?,
170 entity: row.get(2)?,
171 local_id: row.get(3)?,
172 remote_id: row.get(4)?,
173 remote_key: row.get(5)?,
174 remote_updated_at: remote_updated_at.as_deref().and_then(time::parse_ts),
175 base_hash: row.get(7)?,
176 last_synced_at: time::parse_ts(&last_synced_at).unwrap_or_else(time::now_usec),
177 })
178}
179
180#[cfg(test)]
181mod tests {
182 use super::*;
183
184 fn conn() -> Connection {
185 let c = Connection::open_in_memory().unwrap();
186 ensure_table(&c).unwrap();
187 c
188 }
189
190 fn new_link(local_id: i64, remote_id: &str, remote_key: &str) -> NewLink {
191 NewLink {
192 provider: crate::PROVIDER_LINEAR.into(),
193 entity: crate::ENTITY_ISSUE.into(),
194 local_id,
195 remote_id: remote_id.into(),
196 remote_key: remote_key.into(),
197 remote_updated_at: Some(time::now_usec()),
198 base_hash: Some("abc123".into()),
199 }
200 }
201
202 #[test]
203 fn ensure_table_is_idempotent() {
204 let c = conn();
205 ensure_table(&c).unwrap();
208 ensure_table(&c).unwrap();
209 }
210
211 #[test]
212 fn upsert_then_read_back_by_either_side() {
213 let c = conn();
214 let saved = upsert(&c, new_link(7, "uuid-1", "ENG-412")).unwrap();
215 assert_eq!(saved.local_id, 7);
216 assert_eq!(saved.remote_key, "ENG-412");
217
218 let by_l = by_local(&c, crate::PROVIDER_LINEAR, crate::ENTITY_ISSUE, 7)
219 .unwrap()
220 .unwrap();
221 let by_r = by_remote(&c, crate::PROVIDER_LINEAR, crate::ENTITY_ISSUE, "uuid-1")
222 .unwrap()
223 .unwrap();
224 assert_eq!(by_l.id, by_r.id);
225 assert_eq!(by_l.base_hash.as_deref(), Some("abc123"));
226 }
227
228 #[test]
229 fn upsert_refreshes_rather_than_duplicating() {
230 let c = conn();
231 upsert(&c, new_link(7, "uuid-1", "ENG-412")).unwrap();
232 let again = upsert(&c, new_link(7, "uuid-1", "PLAT-9")).unwrap();
234 assert_eq!(again.remote_key, "PLAT-9");
235 assert_eq!(list(&c, crate::PROVIDER_LINEAR).unwrap().len(), 1);
236 }
237
238 #[test]
239 fn two_local_issues_cannot_claim_one_remote_issue() {
240 let c = conn();
241 upsert(&c, new_link(7, "uuid-1", "ENG-412")).unwrap();
242 assert!(upsert(&c, new_link(8, "uuid-1", "ENG-412")).is_err());
245 }
246
247 #[test]
248 fn missing_link_reads_as_none_not_error() {
249 let c = conn();
250 assert!(
251 by_local(&c, crate::PROVIDER_LINEAR, crate::ENTITY_ISSUE, 99)
252 .unwrap()
253 .is_none()
254 );
255 }
256}