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, origin, \
32 progress_comment_id";
33
34#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
38pub enum Origin {
39 #[default]
42 Imported,
43 Pushed,
46}
47
48impl Origin {
49 pub fn as_str(self) -> &'static str {
50 match self {
51 Origin::Imported => "imported",
52 Origin::Pushed => "pushed",
53 }
54 }
55
56 pub fn from_db(s: &str) -> Origin {
60 match s {
61 "pushed" => Origin::Pushed,
62 _ => Origin::Imported,
63 }
64 }
65}
66
67#[derive(Debug, Clone)]
69pub struct RemoteLink {
70 pub id: i64,
71 pub provider: String,
72 pub entity: String,
73 pub local_id: i64,
74 pub remote_id: String,
78 pub remote_key: String,
80 pub remote_updated_at: Option<DateTime<Utc>>,
83 pub base_hash: Option<String>,
86 pub last_synced_at: DateTime<Utc>,
87 pub origin: Origin,
89 pub progress_comment_id: Option<String>,
93}
94
95#[derive(Debug, Clone)]
97pub struct NewLink {
98 pub provider: String,
99 pub entity: String,
100 pub local_id: i64,
101 pub remote_id: String,
102 pub remote_key: String,
103 pub remote_updated_at: Option<DateTime<Utc>>,
104 pub base_hash: Option<String>,
105 pub origin: Origin,
108}
109
110pub fn ensure_table(conn: &Connection) -> Result<()> {
113 for stmt in DDL {
114 conn.execute_batch(stmt)?;
115 }
116 if !has_column(conn, "remote_links", "origin")? {
124 conn.execute_batch(
125 "ALTER TABLE remote_links ADD COLUMN origin TEXT NOT NULL DEFAULT 'imported'",
126 )?;
127 }
128 if !has_column(conn, "remote_links", "progress_comment_id")? {
132 conn.execute_batch("ALTER TABLE remote_links ADD COLUMN progress_comment_id TEXT")?;
133 }
134 Ok(())
135}
136
137fn has_column(conn: &Connection, table: &str, column: &str) -> Result<bool> {
139 let mut stmt = conn.prepare("SELECT 1 FROM pragma_table_info(?1) WHERE name = ?2")?;
140 Ok(stmt.exists(params![table, column])?)
141}
142
143pub fn by_local(
145 conn: &Connection,
146 provider: &str,
147 entity: &str,
148 local_id: i64,
149) -> Result<Option<RemoteLink>> {
150 let sql = format!(
151 "SELECT {COLS} FROM remote_links \
152 WHERE provider = ?1 AND entity = ?2 AND local_id = ?3"
153 );
154 let mut stmt = conn.prepare(&sql)?;
155 let mut rows = stmt.query(params![provider, entity, local_id])?;
156 match rows.next()? {
157 Some(row) => Ok(Some(read(row)?)),
158 None => Ok(None),
159 }
160}
161
162pub fn by_remote(
165 conn: &Connection,
166 provider: &str,
167 entity: &str,
168 remote_id: &str,
169) -> Result<Option<RemoteLink>> {
170 let sql = format!(
171 "SELECT {COLS} FROM remote_links \
172 WHERE provider = ?1 AND entity = ?2 AND remote_id = ?3"
173 );
174 let mut stmt = conn.prepare(&sql)?;
175 let mut rows = stmt.query(params![provider, entity, remote_id])?;
176 match rows.next()? {
177 Some(row) => Ok(Some(read(row)?)),
178 None => Ok(None),
179 }
180}
181
182pub fn upsert(conn: &Connection, new: NewLink) -> Result<RemoteLink> {
189 let now = time::format_usec(time::now_usec());
190 let remote_updated = new.remote_updated_at.map(time::format_usec);
191
192 conn.execute(
196 "INSERT INTO remote_links \
197 (provider, entity, local_id, remote_id, remote_key, \
198 remote_updated_at, base_hash, last_synced_at, inserted_at, updated_at, origin) \
199 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?8, ?8, ?9) \
200 ON CONFLICT (provider, entity, local_id) DO UPDATE SET \
201 remote_id = excluded.remote_id, \
202 remote_key = excluded.remote_key, \
203 remote_updated_at = excluded.remote_updated_at, \
204 base_hash = excluded.base_hash, \
205 last_synced_at = excluded.last_synced_at, \
206 updated_at = excluded.updated_at",
207 params![
208 new.provider,
209 new.entity,
210 new.local_id,
211 new.remote_id,
212 new.remote_key,
213 remote_updated,
214 new.base_hash,
215 now,
216 new.origin.as_str(),
217 ],
218 )?;
219
220 by_local(conn, &new.provider, &new.entity, new.local_id)?.ok_or(Error::NotFound)
221}
222
223pub fn set_progress_comment(
229 conn: &Connection,
230 provider: &str,
231 entity: &str,
232 local_id: i64,
233 comment_id: Option<&str>,
234) -> Result<()> {
235 let n = conn.execute(
236 "UPDATE remote_links SET progress_comment_id = ?4, updated_at = ?5 \
237 WHERE provider = ?1 AND entity = ?2 AND local_id = ?3",
238 params![
239 provider,
240 entity,
241 local_id,
242 comment_id,
243 time::format_usec(time::now_usec()),
244 ],
245 )?;
246 if n == 0 {
247 return Err(Error::NotFound);
248 }
249 Ok(())
250}
251
252pub fn list(conn: &Connection, provider: &str) -> Result<Vec<RemoteLink>> {
254 let sql = format!("SELECT {COLS} FROM remote_links WHERE provider = ?1 ORDER BY id");
255 let mut stmt = conn.prepare(&sql)?;
256 let mut out = Vec::new();
257 let mut rows = stmt.query(params![provider])?;
258 while let Some(row) = rows.next()? {
259 out.push(read(row)?);
260 }
261 Ok(out)
262}
263
264fn read(row: &Row) -> Result<RemoteLink> {
265 let remote_updated_at: Option<String> = row.get(6)?;
266 let last_synced_at: String = row.get(8)?;
267 let origin: String = row.get(9)?;
268 Ok(RemoteLink {
269 id: row.get(0)?,
270 provider: row.get(1)?,
271 entity: row.get(2)?,
272 local_id: row.get(3)?,
273 remote_id: row.get(4)?,
274 remote_key: row.get(5)?,
275 remote_updated_at: remote_updated_at.as_deref().and_then(time::parse_ts),
276 base_hash: row.get(7)?,
277 last_synced_at: time::parse_ts(&last_synced_at).unwrap_or_else(time::now_usec),
278 origin: Origin::from_db(&origin),
279 progress_comment_id: row.get(10)?,
280 })
281}
282
283#[cfg(test)]
284mod tests {
285 use super::*;
286
287 fn conn() -> Connection {
288 let c = Connection::open_in_memory().unwrap();
289 ensure_table(&c).unwrap();
290 c
291 }
292
293 fn new_link(local_id: i64, remote_id: &str, remote_key: &str) -> NewLink {
294 NewLink {
295 provider: crate::PROVIDER_LINEAR.into(),
296 entity: crate::ENTITY_ISSUE.into(),
297 local_id,
298 remote_id: remote_id.into(),
299 remote_key: remote_key.into(),
300 remote_updated_at: Some(time::now_usec()),
301 base_hash: Some("abc123".into()),
302 origin: Origin::Imported,
303 }
304 }
305
306 #[test]
307 fn ensure_table_is_idempotent() {
308 let c = conn();
309 ensure_table(&c).unwrap();
312 ensure_table(&c).unwrap();
313 }
314
315 #[test]
316 fn upsert_then_read_back_by_either_side() {
317 let c = conn();
318 let saved = upsert(&c, new_link(7, "uuid-1", "ENG-412")).unwrap();
319 assert_eq!(saved.local_id, 7);
320 assert_eq!(saved.remote_key, "ENG-412");
321
322 let by_l = by_local(&c, crate::PROVIDER_LINEAR, crate::ENTITY_ISSUE, 7)
323 .unwrap()
324 .unwrap();
325 let by_r = by_remote(&c, crate::PROVIDER_LINEAR, crate::ENTITY_ISSUE, "uuid-1")
326 .unwrap()
327 .unwrap();
328 assert_eq!(by_l.id, by_r.id);
329 assert_eq!(by_l.base_hash.as_deref(), Some("abc123"));
330 }
331
332 #[test]
333 fn upsert_refreshes_rather_than_duplicating() {
334 let c = conn();
335 upsert(&c, new_link(7, "uuid-1", "ENG-412")).unwrap();
336 let again = upsert(&c, new_link(7, "uuid-1", "PLAT-9")).unwrap();
338 assert_eq!(again.remote_key, "PLAT-9");
339 assert_eq!(list(&c, crate::PROVIDER_LINEAR).unwrap().len(), 1);
340 }
341
342 #[test]
343 fn two_local_issues_cannot_claim_one_remote_issue() {
344 let c = conn();
345 upsert(&c, new_link(7, "uuid-1", "ENG-412")).unwrap();
346 assert!(upsert(&c, new_link(8, "uuid-1", "ENG-412")).is_err());
349 }
350
351 #[test]
352 fn missing_link_reads_as_none_not_error() {
353 let c = conn();
354 assert!(
355 by_local(&c, crate::PROVIDER_LINEAR, crate::ENTITY_ISSUE, 99)
356 .unwrap()
357 .is_none()
358 );
359 }
360
361 #[test]
362 fn ensure_table_adds_origin_to_a_pre_origin_table_and_backfills_imported() {
363 let c = Connection::open_in_memory().unwrap();
364 c.execute_batch(
367 r#"CREATE TABLE "remote_links" (
368 "id" INTEGER PRIMARY KEY AUTOINCREMENT,
369 "provider" TEXT NOT NULL,
370 "entity" TEXT NOT NULL,
371 "local_id" INTEGER NOT NULL,
372 "remote_id" TEXT NOT NULL,
373 "remote_key" TEXT NOT NULL,
374 "remote_updated_at" TEXT,
375 "base_hash" TEXT,
376 "last_synced_at" TEXT NOT NULL,
377 "inserted_at" TEXT NOT NULL,
378 "updated_at" TEXT NOT NULL
379 )"#,
380 )
381 .unwrap();
382 c.execute(
383 "INSERT INTO remote_links \
384 (provider, entity, local_id, remote_id, remote_key, last_synced_at, \
385 inserted_at, updated_at) \
386 VALUES ('linear', 'issue', 7, 'uuid-1', 'ENG-412', '2026-01-01T00:00:00Z', \
387 '2026-01-01T00:00:00Z', '2026-01-01T00:00:00Z')",
388 [],
389 )
390 .unwrap();
391
392 ensure_table(&c).unwrap();
393 ensure_table(&c).unwrap();
395
396 let link = by_local(&c, crate::PROVIDER_LINEAR, crate::ENTITY_ISSUE, 7)
397 .unwrap()
398 .unwrap();
399 assert_eq!(
400 link.origin,
401 Origin::Imported,
402 "pre-origin rows were only ever created by import"
403 );
404 }
405
406 #[test]
407 fn upsert_records_origin_and_a_later_sync_never_flips_it() {
408 let c = conn();
409 let mut first = new_link(7, "uuid-1", "ENG-412");
410 first.origin = Origin::Pushed;
411 let saved = upsert(&c, first).unwrap();
412 assert_eq!(saved.origin, Origin::Pushed);
413
414 let again = upsert(&c, new_link(7, "uuid-1", "ENG-412")).unwrap();
417 assert_eq!(again.origin, Origin::Pushed);
418 }
419
420 #[test]
421 fn ensure_table_adds_progress_comment_id_to_a_pre_column_table() {
422 let c = Connection::open_in_memory().unwrap();
423 c.execute_batch(
425 r#"CREATE TABLE "remote_links" (
426 "id" INTEGER PRIMARY KEY AUTOINCREMENT,
427 "provider" TEXT NOT NULL,
428 "entity" TEXT NOT NULL,
429 "local_id" INTEGER NOT NULL,
430 "remote_id" TEXT NOT NULL,
431 "remote_key" TEXT NOT NULL,
432 "remote_updated_at" TEXT,
433 "base_hash" TEXT,
434 "last_synced_at" TEXT NOT NULL,
435 "inserted_at" TEXT NOT NULL,
436 "updated_at" TEXT NOT NULL,
437 "origin" TEXT NOT NULL DEFAULT 'imported'
438 )"#,
439 )
440 .unwrap();
441 c.execute(
442 "INSERT INTO remote_links \
443 (provider, entity, local_id, remote_id, remote_key, last_synced_at, \
444 inserted_at, updated_at) \
445 VALUES ('linear', 'issue', 7, 'uuid-1', 'ENG-412', '2026-01-01T00:00:00Z', \
446 '2026-01-01T00:00:00Z', '2026-01-01T00:00:00Z')",
447 [],
448 )
449 .unwrap();
450
451 ensure_table(&c).unwrap();
452 ensure_table(&c).unwrap();
454
455 let link = by_local(&c, crate::PROVIDER_LINEAR, crate::ENTITY_ISSUE, 7)
456 .unwrap()
457 .unwrap();
458 assert_eq!(
459 link.progress_comment_id, None,
460 "a pre-column row has no living comment yet"
461 );
462 }
463
464 #[test]
465 fn set_progress_comment_persists_and_reads_back() {
466 let c = conn();
467 upsert(&c, new_link(7, "uuid-1", "ENG-412")).unwrap();
468 set_progress_comment(
469 &c,
470 crate::PROVIDER_LINEAR,
471 crate::ENTITY_ISSUE,
472 7,
473 Some("comment-uuid-1"),
474 )
475 .unwrap();
476 let link = by_local(&c, crate::PROVIDER_LINEAR, crate::ENTITY_ISSUE, 7)
477 .unwrap()
478 .unwrap();
479 assert_eq!(link.progress_comment_id.as_deref(), Some("comment-uuid-1"));
480
481 set_progress_comment(&c, crate::PROVIDER_LINEAR, crate::ENTITY_ISSUE, 7, None).unwrap();
483 let link = by_local(&c, crate::PROVIDER_LINEAR, crate::ENTITY_ISSUE, 7)
484 .unwrap()
485 .unwrap();
486 assert_eq!(link.progress_comment_id, None);
487 }
488
489 #[test]
490 fn upsert_never_clobbers_a_stored_comment_id() {
491 let c = conn();
492 upsert(&c, new_link(7, "uuid-1", "ENG-412")).unwrap();
493 set_progress_comment(
494 &c,
495 crate::PROVIDER_LINEAR,
496 crate::ENTITY_ISSUE,
497 7,
498 Some("comment-uuid-1"),
499 )
500 .unwrap();
501
502 let again = upsert(&c, new_link(7, "uuid-1", "ENG-412")).unwrap();
505 assert_eq!(again.progress_comment_id.as_deref(), Some("comment-uuid-1"));
506 }
507
508 #[test]
509 fn an_unrecognized_origin_reads_as_imported() {
510 assert_eq!(Origin::from_db("imported"), Origin::Imported);
513 assert_eq!(Origin::from_db("pushed"), Origin::Pushed);
514 assert_eq!(Origin::from_db("mystery"), Origin::Imported);
515 }
516}