1mod hydration;
2mod projects;
3pub mod schema;
4mod sync;
5pub use hydration::*;
6pub use projects::*;
7pub use sync::*;
8#[cfg(test)]
9pub(crate) mod test_helpers;
10
11use anyhow::{Context, Result};
12use rusqlite::Connection;
13use serde::{Deserialize, Serialize};
14use std::path::Path;
15use std::sync::{Arc, Mutex};
16
17pub struct BlockerRow {
18 pub issue_id: String,
19 pub identifier: String,
20 pub title: String,
21 pub state_name: String,
22 pub state_type: String,
23}
24
25#[derive(Clone)]
26pub struct Database {
27 conn: Arc<Mutex<Connection>>,
28}
29
30impl Database {
31 pub fn open(path: &Path) -> Result<Self> {
32 let conn = Connection::open(path)
33 .with_context(|| format!("Failed to open database at {}", path.display()))?;
34
35 conn.execute_batch("PRAGMA journal_mode=WAL; PRAGMA foreign_keys=ON;")?;
36
37 let db = Self {
38 conn: Arc::new(Mutex::new(conn)),
39 };
40 db.migrate()?;
41 Ok(db)
42 }
43
44 fn migrate(&self) -> Result<()> {
45 let conn = self.conn.lock().unwrap();
46 schema::run_migrations(&conn)?;
47 conn.execute(
51 "UPDATE issue_hydration_state
52 SET status='retryable', next_retry_at=datetime('now'),
53 queue_reason='retry'
54 WHERE status='running'",
55 [],
56 )?;
57 Ok(())
58 }
59
60 pub fn with_conn<F, T>(&self, f: F) -> Result<T>
61 where
62 F: FnOnce(&Connection) -> Result<T>,
63 {
64 let conn = self.conn.lock().unwrap();
65 f(&conn)
66 }
67
68 pub fn upsert_workspace(
71 &self,
72 id: &str,
73 linear_org_id: Option<&str>,
74 display_name: Option<&str>,
75 ) -> Result<()> {
76 self.with_conn(|conn| {
77 conn.execute(
78 "INSERT INTO workspaces (id, linear_org_id, display_name)
79 VALUES (?1, ?2, ?3)
80 ON CONFLICT(id) DO UPDATE SET
81 linear_org_id=excluded.linear_org_id,
82 display_name=excluded.display_name",
83 rusqlite::params![id, linear_org_id, display_name],
84 )?;
85 Ok(())
86 })
87 }
88
89 pub fn get_workspace(&self, id: &str) -> Result<Option<WorkspaceRow>> {
90 self.with_conn(|conn| {
91 let mut stmt = conn.prepare(
92 "SELECT id, linear_org_id, display_name, created_at FROM workspaces WHERE id = ?1",
93 )?;
94 let mut rows = stmt.query(rusqlite::params![id])?;
95 if let Some(row) = rows.next()? {
96 Ok(Some(WorkspaceRow {
97 id: row.get(0)?,
98 linear_org_id: row.get(1)?,
99 display_name: row.get(2)?,
100 created_at: row.get(3)?,
101 }))
102 } else {
103 Ok(None)
104 }
105 })
106 }
107
108 pub fn list_workspaces(&self) -> Result<Vec<WorkspaceRow>> {
109 self.with_conn(|conn| {
110 let mut stmt = conn.prepare(
111 "SELECT id, linear_org_id, display_name, created_at FROM workspaces ORDER BY id",
112 )?;
113 let rows = stmt.query_map([], |row| {
114 Ok(WorkspaceRow {
115 id: row.get(0)?,
116 linear_org_id: row.get(1)?,
117 display_name: row.get(2)?,
118 created_at: row.get(3)?,
119 })
120 })?;
121 Ok(rows.collect::<std::result::Result<Vec<_>, _>>()?)
122 })
123 }
124
125 pub fn delete_workspace(&self, id: &str) -> Result<usize> {
127 self.with_conn(|conn| {
128 let issue_count: usize = conn.query_row(
130 "SELECT COUNT(*) FROM issues WHERE workspace_id = ?1",
131 rusqlite::params![id],
132 |row| row.get(0),
133 )?;
134 conn.execute(
135 "DELETE FROM issues WHERE workspace_id = ?1",
136 rusqlite::params![id],
137 )?;
138 conn.execute(
139 "DELETE FROM comments WHERE workspace_id = ?1",
140 rusqlite::params![id],
141 )?;
142 conn.execute(
143 "DELETE FROM comment_sync_state WHERE workspace_id = ?1",
144 rusqlite::params![id],
145 )?;
146 conn.execute(
147 "DELETE FROM sync_state WHERE workspace_id = ?1",
148 rusqlite::params![id],
149 )?;
150 conn.execute(
151 "DELETE FROM sync_family_state WHERE workspace_id = ?1",
152 rusqlite::params![id],
153 )?;
154 conn.execute(
155 "DELETE FROM cycles WHERE workspace_id = ?1",
156 rusqlite::params![id],
157 )?;
158 conn.execute(
159 "DELETE FROM labels WHERE workspace_id = ?1",
160 rusqlite::params![id],
161 )?;
162 conn.execute(
163 "DELETE FROM projects WHERE workspace_id = ?1",
164 rusqlite::params![id],
165 )?;
166 conn.execute(
167 "DELETE FROM workspaces WHERE id = ?1",
168 rusqlite::params![id],
169 )?;
170 Ok(issue_count)
171 })
172 }
173
174 pub fn upsert_label(&self, label: &Label) -> Result<()> {
177 self.with_conn(|conn| {
178 conn.execute(
179 "INSERT INTO labels (id, workspace_id, name, color, parent_id)
180 VALUES (?1, ?2, ?3, ?4, ?5)
181 ON CONFLICT(id) DO UPDATE SET
182 workspace_id=excluded.workspace_id,
183 name=excluded.name,
184 color=excluded.color,
185 parent_id=excluded.parent_id",
186 rusqlite::params![
187 label.id,
188 label.workspace_id,
189 label.name,
190 label.color,
191 label.parent_id
192 ],
193 )?;
194 Ok(())
195 })
196 }
197
198 pub fn list_labels(&self, workspace_id: &str) -> Result<Vec<Label>> {
199 self.with_conn(|conn| {
200 let mut stmt = conn.prepare(
201 "SELECT id, workspace_id, name, color, parent_id
202 FROM labels WHERE workspace_id = ?1
203 ORDER BY name COLLATE NOCASE ASC",
204 )?;
205 let rows = stmt.query_map(rusqlite::params![workspace_id], |row| {
206 Ok(Label {
207 id: row.get(0)?,
208 workspace_id: row.get(1)?,
209 name: row.get(2)?,
210 color: row.get(3)?,
211 parent_id: row.get(4)?,
212 })
213 })?;
214 Ok(rows.collect::<std::result::Result<Vec<_>, _>>()?)
215 })
216 }
217
218 pub fn delete_labels_for_workspace_not_in(
221 &self,
222 workspace_id: &str,
223 keep_ids: &[String],
224 ) -> Result<usize> {
225 self.with_conn(|conn| {
226 if keep_ids.is_empty() {
227 let n = conn.execute(
228 "DELETE FROM labels WHERE workspace_id = ?1",
229 rusqlite::params![workspace_id],
230 )?;
231 return Ok(n);
232 }
233 let placeholders = (0..keep_ids.len())
234 .map(|i| format!("?{}", i + 2))
235 .collect::<Vec<_>>()
236 .join(", ");
237 let sql = format!(
238 "DELETE FROM labels WHERE workspace_id = ?1 AND id NOT IN ({placeholders})"
239 );
240 let mut params: Vec<Box<dyn rusqlite::types::ToSql>> =
241 vec![Box::new(workspace_id.to_string())];
242 for id in keep_ids {
243 params.push(Box::new(id.clone()));
244 }
245 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
246 params.iter().map(|p| p.as_ref()).collect();
247 let n = conn.execute(&sql, param_refs.as_slice())?;
248 Ok(n)
249 })
250 }
251
252 pub fn resolve_label_ids_local(
255 &self,
256 workspace_id: &str,
257 names: &[String],
258 ) -> Result<(Vec<String>, Vec<String>)> {
259 if names.is_empty() {
260 return Ok((Vec::new(), Vec::new()));
261 }
262 self.with_conn(|conn| {
263 let mut resolved = Vec::new();
264 let mut unknown = Vec::new();
265 let mut stmt = conn.prepare(
266 "SELECT id FROM labels WHERE workspace_id = ?1 AND name = ?2 COLLATE NOCASE",
267 )?;
268 for name in names {
269 let mut rows = stmt.query(rusqlite::params![workspace_id, name])?;
270 if let Some(row) = rows.next()? {
271 resolved.push(row.get::<_, String>(0)?);
272 } else {
273 unknown.push(name.clone());
274 }
275 }
276 Ok((resolved, unknown))
277 })
278 }
279
280 pub fn replace_issue_labels(&self, issue_id: &str, label_ids: &[String]) -> Result<()> {
285 self.with_conn(|conn| {
286 let tx = conn.unchecked_transaction()?;
287 tx.execute(
288 "DELETE FROM issue_labels WHERE issue_id = ?1",
289 rusqlite::params![issue_id],
290 )?;
291 for lid in label_ids {
292 let exists: i64 = tx.query_row(
293 "SELECT COUNT(*) FROM labels WHERE id = ?1",
294 rusqlite::params![lid],
295 |r| r.get(0),
296 )?;
297 if exists == 0 {
298 eprintln!(
299 "warning: skipping unknown label id '{}' for issue '{}'",
300 lid, issue_id
301 );
302 continue;
303 }
304 tx.execute(
305 "INSERT OR IGNORE INTO issue_labels (issue_id, label_id) VALUES (?1, ?2)",
306 rusqlite::params![issue_id, lid],
307 )?;
308 }
309 tx.commit()?;
310 Ok(())
311 })
312 }
313
314 pub fn get_issue_label_ids(&self, issue_id: &str) -> Result<Vec<String>> {
315 self.with_conn(|conn| {
316 let mut stmt = conn.prepare(
317 "SELECT label_id FROM issue_labels WHERE issue_id = ?1 ORDER BY label_id",
318 )?;
319 let rows =
320 stmt.query_map(rusqlite::params![issue_id], |row| row.get::<_, String>(0))?;
321 Ok(rows.collect::<std::result::Result<Vec<_>, _>>()?)
322 })
323 }
324
325 pub fn upsert_issue(&self, issue: &Issue) -> Result<()> {
328 self.upsert_issue_with_label_policy(issue, false)
329 }
330
331 pub fn upsert_issue_preserving_labels(&self, issue: &Issue) -> Result<()> {
332 self.upsert_issue_with_label_policy(issue, true)
333 }
334
335 fn upsert_issue_with_label_policy(&self, issue: &Issue, preserve_labels: bool) -> Result<()> {
336 let embedding_content_hash =
337 crate::embedding::issue_content_hash(&issue.title, issue.description.as_deref());
338 self.with_conn(|conn| {
339 conn.execute(
340 "INSERT INTO issues (id, identifier, team_key, title, description, state_name, state_type, priority, assignee_name, project_name, labels_json, created_at, updated_at, content_hash, synced_at, url, branch_name, workspace_id, project_id, project_milestone_id, project_milestone_name, cycle_id, cycle_name, archived_at, embedding_content_hash)
341 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, datetime('now'), ?15, ?16, ?17, ?18, ?19, ?20, ?21, ?22, ?23, ?24)
342 ON CONFLICT(id) DO UPDATE SET
343 identifier=excluded.identifier, team_key=excluded.team_key, title=excluded.title,
344 description=excluded.description, state_name=excluded.state_name, state_type=excluded.state_type,
345 priority=excluded.priority, assignee_name=excluded.assignee_name, project_name=excluded.project_name,
346 labels_json=CASE WHEN ?25 THEN issues.labels_json ELSE excluded.labels_json END,
347 updated_at=excluded.updated_at,
348 content_hash=CASE WHEN ?25 THEN issues.content_hash ELSE excluded.content_hash END,
349 embedding_content_hash=excluded.embedding_content_hash,
350 url=excluded.url, branch_name=excluded.branch_name,
351 workspace_id=excluded.workspace_id, project_id=excluded.project_id,
352 project_milestone_id=excluded.project_milestone_id,
353 project_milestone_name=excluded.project_milestone_name,
354 cycle_id=excluded.cycle_id,
355 cycle_name=excluded.cycle_name,
356 archived_at=excluded.archived_at,
357 synced_at=datetime('now')",
358 rusqlite::params![
359 issue.id, issue.identifier, issue.team_key, issue.title, issue.description,
360 issue.state_name, issue.state_type, issue.priority, issue.assignee_name,
361 issue.project_name, issue.labels_json, issue.created_at, issue.updated_at,
362 issue.content_hash, issue.url, issue.branch_name, issue.workspace_id,
363 issue.project_id, issue.project_milestone_id, issue.project_milestone_name,
364 issue.cycle_id, issue.cycle_name,
365 issue.archived_at,
366 embedding_content_hash,
367 preserve_labels,
368 ],
369 )?;
370 Ok(())
371 })
372 }
373
374 pub fn get_issue(&self, id_or_identifier: &str) -> Result<Option<Issue>> {
375 self.with_conn(|conn| {
376 let mut stmt = conn.prepare(
377 "SELECT id, identifier, team_key, title, description, state_name, state_type, priority, assignee_name, project_name, labels_json, created_at, updated_at, content_hash, synced_at, url, branch_name, workspace_id, project_id, project_milestone_id, project_milestone_name, cycle_id, cycle_name, archived_at
378 FROM issues WHERE id = ?1 OR identifier = ?1"
379 )?;
380 let mut rows = stmt.query(rusqlite::params![id_or_identifier])?;
381 if let Some(row) = rows.next()? {
382 Ok(Some(Issue::from_row(row)?))
383 } else {
384 Ok(None)
385 }
386 })
387 }
388
389 fn label_filter_fragment(
395 label_ids: &[String],
396 param_offset: usize,
397 table_alias: &str,
398 ) -> (String, Vec<Box<dyn rusqlite::types::ToSql>>) {
399 let n = label_ids.len();
400 let placeholders = (0..n)
401 .map(|i| format!("?{}", param_offset + i))
402 .collect::<Vec<_>>()
403 .join(", ");
404 let sql = format!(
405 "{table_alias}.id IN (\
406 SELECT issue_id FROM issue_labels \
407 WHERE label_id IN ({placeholders}) \
408 GROUP BY issue_id \
409 HAVING COUNT(DISTINCT label_id) = {n}\
410 )"
411 );
412 let params: Vec<Box<dyn rusqlite::types::ToSql>> = label_ids
413 .iter()
414 .map(|s| Box::new(s.clone()) as Box<dyn rusqlite::types::ToSql>)
415 .collect();
416 (sql, params)
417 }
418
419 pub fn get_unprioritized_issues(
420 &self,
421 team_key: Option<&str>,
422 include_completed: bool,
423 workspace_id: &str,
424 ) -> Result<Vec<Issue>> {
425 self.get_unprioritized_issues_filtered(team_key, include_completed, workspace_id, None)
426 }
427
428 pub fn get_unprioritized_issues_filtered(
429 &self,
430 team_key: Option<&str>,
431 include_completed: bool,
432 workspace_id: &str,
433 label_ids: Option<&[String]>,
434 ) -> Result<Vec<Issue>> {
435 self.with_conn(|conn| {
436 let state_filter = if include_completed {
437 ""
438 } else {
439 " AND state_type NOT IN ('completed', 'canceled')"
440 };
441
442 let mut params: Vec<Box<dyn rusqlite::types::ToSql>> = Vec::new();
444 let base_where: String = if let Some(team) = team_key {
445 params.push(Box::new(team.to_string()));
446 params.push(Box::new(workspace_id.to_string()));
447 "team_key = ?1 AND workspace_id = ?2".to_string()
448 } else {
449 params.push(Box::new(workspace_id.to_string()));
450 "workspace_id = ?1".to_string()
451 };
452
453 let label_clause = if let Some(ids) = label_ids.filter(|ids| !ids.is_empty()) {
454 let (frag, mut lp) = Self::label_filter_fragment(ids, params.len() + 1, "issues");
455 params.append(&mut lp);
456 format!(" AND {frag}")
457 } else {
458 String::new()
459 };
460
461 let sql = format!(
462 "SELECT id, identifier, team_key, title, description, state_name, state_type, priority, assignee_name, project_name, labels_json, created_at, updated_at, content_hash, synced_at, url, branch_name, workspace_id, project_id, project_milestone_id, project_milestone_name, cycle_id, cycle_name
463 FROM issues WHERE priority = 0{state_filter} AND {base_where}{label_clause}
464 ORDER BY created_at DESC"
465 );
466
467 let mut stmt = conn.prepare(&sql)?;
468 let param_refs: Vec<&dyn rusqlite::types::ToSql> = params.iter().map(|p| p.as_ref()).collect();
469 let rows = stmt.query_map(param_refs.as_slice(), |row| Ok(Issue::from_row(row).unwrap()))?;
470 Ok(rows.collect::<std::result::Result<Vec<_>, _>>()?)
471 })
472 }
473
474 pub fn get_issues_by_state_types(
475 &self,
476 team_key: &str,
477 state_types: &[String],
478 workspace_id: &str,
479 ) -> Result<Vec<Issue>> {
480 self.with_conn(|conn| {
481 let placeholders: String = state_types
482 .iter()
483 .enumerate()
484 .map(|(i, _)| format!("?{}", i + 3))
485 .collect::<Vec<_>>()
486 .join(", ");
487 let sql = format!(
488 "SELECT id, identifier, team_key, title, description, state_name, state_type, \
489 priority, assignee_name, project_name, labels_json, created_at, updated_at, \
490 content_hash, synced_at, url, branch_name, workspace_id, project_id, \
491 project_milestone_id, project_milestone_name, cycle_id, cycle_name \
492 FROM issues WHERE team_key = ?1 AND workspace_id = ?2 AND state_type IN ({placeholders}) \
493 ORDER BY priority ASC, created_at DESC"
494 );
495 let mut stmt = conn.prepare(&sql)?;
496 let mut params: Vec<Box<dyn rusqlite::types::ToSql>> =
497 vec![Box::new(team_key.to_string()), Box::new(workspace_id.to_string())];
498 for st in state_types {
499 params.push(Box::new(st.clone()));
500 }
501 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
502 params.iter().map(|p| p.as_ref()).collect();
503 let rows = stmt.query_map(param_refs.as_slice(), |row| {
504 Ok(Issue::from_row(row).unwrap())
505 })?;
506 Ok(rows.collect::<std::result::Result<Vec<_>, _>>()?)
507 })
508 }
509
510 pub fn get_blockers_for_issues(&self, issue_ids: &[String]) -> Result<Vec<BlockerRow>> {
513 if issue_ids.is_empty() {
514 return Ok(vec![]);
515 }
516 self.with_conn(|conn| {
517 let placeholders: String = issue_ids
518 .iter()
519 .enumerate()
520 .map(|(i, _)| format!("?{}", i + 1))
521 .collect::<Vec<_>>()
522 .join(", ");
523
524 let sql_fwd = format!(
526 "SELECT r.issue_id, COALESCE(i.identifier, r.related_issue_identifier),
527 COALESCE(i.title, ''), COALESCE(i.state_name, ''), COALESCE(i.state_type, '')
528 FROM issue_relations r
529 LEFT JOIN issues i ON r.related_issue_id = i.id
530 WHERE r.issue_id IN ({placeholders}) AND r.relation_type = 'blocked_by'"
531 );
532
533 let sql_inv = format!(
535 "SELECT r.related_issue_id, i2.identifier,
536 COALESCE(i2.title, ''), COALESCE(i2.state_name, ''), COALESCE(i2.state_type, '')
537 FROM issue_relations r
538 JOIN issues i ON r.related_issue_id = i.id
539 JOIN issues i2 ON r.issue_id = i2.id
540 WHERE r.related_issue_id IN ({placeholders}) AND r.relation_type = 'blocks'"
541 );
542
543 let mut results = Vec::new();
544 let params: Vec<Box<dyn rusqlite::types::ToSql>> =
545 issue_ids.iter().map(|id| Box::new(id.clone()) as _).collect();
546 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
547 params.iter().map(|p| p.as_ref()).collect();
548
549 for sql in [&sql_fwd, &sql_inv] {
550 let mut stmt = conn.prepare(sql)?;
551 let rows = stmt.query_map(param_refs.as_slice(), |row| {
552 Ok(BlockerRow {
553 issue_id: row.get(0)?,
554 identifier: row.get(1)?,
555 title: row.get(2)?,
556 state_name: row.get(3)?,
557 state_type: row.get(4)?,
558 })
559 })?;
560 for row in rows {
561 results.push(row?);
562 }
563 }
564 Ok(results)
565 })
566 }
567
568 pub fn count_issues(&self, team_key: Option<&str>, workspace_id: &str) -> Result<usize> {
569 self.with_conn(|conn| {
570 let count: usize = if let Some(team) = team_key {
571 conn.query_row(
572 "SELECT COUNT(*) FROM issues WHERE team_key = ?1 AND workspace_id = ?2",
573 rusqlite::params![team, workspace_id],
574 |row| row.get(0),
575 )?
576 } else {
577 conn.query_row(
578 "SELECT COUNT(*) FROM issues WHERE workspace_id = ?1",
579 rusqlite::params![workspace_id],
580 |row| row.get(0),
581 )?
582 };
583 Ok(count)
584 })
585 }
586
587 pub fn get_field_completeness(
589 &self,
590 team_key: Option<&str>,
591 workspace_id: &str,
592 ) -> Result<(usize, usize, usize, usize, usize)> {
593 self.with_conn(|conn| {
594 let (sql, params): (String, Vec<Box<dyn rusqlite::types::ToSql>>) =
595 if let Some(team) = team_key {
596 (
597 "SELECT COUNT(*),
598 SUM(CASE WHEN description IS NOT NULL AND description != '' THEN 1 ELSE 0 END),
599 SUM(CASE WHEN priority > 0 THEN 1 ELSE 0 END),
600 SUM(CASE WHEN labels_json != '[]' THEN 1 ELSE 0 END),
601 SUM(CASE WHEN project_name IS NOT NULL AND project_name != '' THEN 1 ELSE 0 END)
602 FROM issues WHERE team_key = ?1 AND workspace_id = ?2"
603 .to_string(),
604 vec![Box::new(team.to_string()) as Box<dyn rusqlite::types::ToSql>, Box::new(workspace_id.to_string())],
605 )
606 } else {
607 (
608 "SELECT COUNT(*),
609 SUM(CASE WHEN description IS NOT NULL AND description != '' THEN 1 ELSE 0 END),
610 SUM(CASE WHEN priority > 0 THEN 1 ELSE 0 END),
611 SUM(CASE WHEN labels_json != '[]' THEN 1 ELSE 0 END),
612 SUM(CASE WHEN project_name IS NOT NULL AND project_name != '' THEN 1 ELSE 0 END)
613 FROM issues WHERE workspace_id = ?1"
614 .to_string(),
615 vec![Box::new(workspace_id.to_string()) as Box<dyn rusqlite::types::ToSql>],
616 )
617 };
618 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
619 params.iter().map(|p| p.as_ref()).collect();
620 let row = conn.query_row(&sql, param_refs.as_slice(), |row| {
621 Ok((
622 row.get::<_, usize>(0)?,
623 row.get::<_, Option<usize>>(1)?.unwrap_or(0),
624 row.get::<_, Option<usize>>(2)?.unwrap_or(0),
625 row.get::<_, Option<usize>>(3)?.unwrap_or(0),
626 row.get::<_, Option<usize>>(4)?.unwrap_or(0),
627 ))
628 })?;
629 Ok(row)
630 })
631 }
632
633 #[allow(unused_assignments)]
636 pub fn list_all_issues(
637 &self,
638 team_key: Option<&str>,
639 filter: Option<&str>,
640 limit: usize,
641 offset: usize,
642 workspace_id: &str,
643 ) -> Result<Vec<IssueSummary>> {
644 self.with_conn(|conn| {
645 let mut conditions = Vec::new();
646 let mut params: Vec<Box<dyn rusqlite::types::ToSql>> = Vec::new();
647 let mut param_idx = 1;
648
649 conditions.push(format!("i.workspace_id = ?{param_idx}"));
651 params.push(Box::new(workspace_id.to_string()));
652 param_idx += 1;
653
654 if let Some(team) = team_key {
655 conditions.push(format!("i.team_key = ?{param_idx}"));
656 params.push(Box::new(team.to_string()));
657 param_idx += 1;
658 }
659
660 if let Some(text) = filter {
661 let like = format!("%{text}%");
662 conditions.push(format!(
663 "(i.identifier LIKE ?{} OR i.title LIKE ?{})",
664 param_idx,
665 param_idx + 1
666 ));
667 params.push(Box::new(like.clone()));
668 params.push(Box::new(like));
669 param_idx += 2;
670 }
671
672 let _ = param_idx;
673
674 let where_clause = if conditions.is_empty() {
675 String::new()
676 } else {
677 format!("WHERE {}", conditions.join(" AND "))
678 };
679
680 let limit_idx = params.len() + 1;
681 let offset_idx = params.len() + 2;
682
683 let sql = format!(
684 "SELECT i.id, i.identifier, i.team_key, i.title, i.state_name, i.state_type,
685 i.priority, i.project_name, i.labels_json, i.updated_at, i.url,
686 i.description IS NOT NULL AND i.description != '' AS has_desc,
687 EXISTS(SELECT 1 FROM chunks c WHERE c.issue_id = i.id) AS has_emb
688 FROM issues i
689 {where_clause}
690 ORDER BY i.updated_at DESC
691 LIMIT ?{limit_idx} OFFSET ?{offset_idx}"
692 );
693 params.push(Box::new(limit as i64));
694 params.push(Box::new(offset as i64));
695
696 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
697 params.iter().map(|p| p.as_ref()).collect();
698 let mut stmt = conn.prepare(&sql)?;
699 let rows = stmt.query_map(param_refs.as_slice(), |row| {
700 let labels_json: String = row.get(8)?;
701 let labels: Vec<String> = serde_json::from_str(&labels_json).unwrap_or_default();
702 Ok(IssueSummary {
703 id: row.get(0)?,
704 identifier: row.get(1)?,
705 team_key: row.get(2)?,
706 title: row.get(3)?,
707 state_name: row.get(4)?,
708 state_type: row.get(5)?,
709 priority: row.get(6)?,
710 project_name: row.get(7)?,
711 labels,
712 updated_at: row.get(9)?,
713 url: row.get(10)?,
714 has_description: row.get(11)?,
715 has_embedding: row.get(12)?,
716 })
717 })?;
718 Ok(rows.collect::<std::result::Result<Vec<_>, _>>()?)
719 })
720 }
721
722 pub fn upsert_relations(&self, issue_id: &str, relations: &[Relation]) -> Result<()> {
725 self.with_conn(|conn| {
726 conn.execute(
727 "DELETE FROM issue_relations WHERE issue_id = ?1",
728 rusqlite::params![issue_id],
729 )?;
730 let mut stmt = conn.prepare(
731 "INSERT OR IGNORE INTO issue_relations (id, issue_id, related_issue_id, related_issue_identifier, relation_type)
732 VALUES (?1, ?2, ?3, ?4, ?5)"
733 )?;
734 for rel in relations {
735 stmt.execute(rusqlite::params![
736 rel.id, rel.issue_id, rel.related_issue_id,
737 rel.related_issue_identifier, rel.relation_type,
738 ])?;
739 }
740 Ok(())
741 })
742 }
743
744 pub fn get_relations_enriched(&self, issue_id: &str) -> Result<Vec<EnrichedRelation>> {
745 self.with_conn(|conn| {
746 let mut stmt = conn.prepare(
748 "SELECT r.id, r.relation_type, r.related_issue_identifier,
749 COALESCE(i.title, ''), COALESCE(i.state_name, ''), COALESCE(i.url, '')
750 FROM issue_relations r
751 LEFT JOIN issues i ON r.related_issue_id = i.id
752 WHERE r.issue_id = ?1",
753 )?;
754 let forward = stmt
755 .query_map(rusqlite::params![issue_id], |row| {
756 Ok(EnrichedRelation {
757 relation_id: row.get(0)?,
758 relation_type: row.get(1)?,
759 issue_identifier: row.get(2)?,
760 issue_title: row.get(3)?,
761 issue_state: row.get(4)?,
762 issue_url: row.get(5)?,
763 })
764 })?
765 .collect::<std::result::Result<Vec<_>, _>>()?;
766
767 let mut stmt2 = conn.prepare(
769 "SELECT r.id, r.relation_type, i2.identifier,
770 COALESCE(i2.title, ''), COALESCE(i2.state_name, ''), COALESCE(i2.url, '')
771 FROM issue_relations r
772 JOIN issues i ON r.related_issue_id = i.id
773 JOIN issues i2 ON r.issue_id = i2.id
774 WHERE r.related_issue_id = i.id AND i.id = ?1",
775 )?;
776 let inverse = stmt2
777 .query_map(rusqlite::params![issue_id], |row| {
778 let raw_type: String = row.get(1)?;
779 let flipped = match raw_type.as_str() {
780 "blocks" => "blocked_by".to_string(),
781 "blocked_by" => "blocks".to_string(),
782 other => other.to_string(), };
784 Ok(EnrichedRelation {
785 relation_id: row.get(0)?,
786 relation_type: flipped,
787 issue_identifier: row.get(2)?,
788 issue_title: row.get(3)?,
789 issue_state: row.get(4)?,
790 issue_url: row.get(5)?,
791 })
792 })?
793 .collect::<std::result::Result<Vec<_>, _>>()?;
794
795 let mut all = forward;
796 all.extend(inverse);
797 Ok(all)
798 })
799 }
800
801 pub fn find_relation_id(
803 &self,
804 issue_id: &str,
805 related_issue_id: &str,
806 relation_type: &str,
807 ) -> Result<Option<String>> {
808 self.with_conn(|conn| {
809 let mut stmt = conn.prepare(
810 "SELECT id FROM issue_relations WHERE issue_id = ?1 AND related_issue_id = ?2 AND relation_type = ?3"
811 )?;
812 let mut rows = stmt.query(rusqlite::params![issue_id, related_issue_id, relation_type])?;
813 if let Some(row) = rows.next()? {
814 Ok(Some(row.get(0)?))
815 } else {
816 Ok(None)
817 }
818 })
819 }
820
821 pub fn upsert_chunks(&self, issue_id: &str, chunks: &[(usize, String, Vec<u8>)]) -> Result<()> {
824 self.upsert_chunks_with_model(issue_id, chunks, "")
825 }
826
827 pub fn upsert_chunks_with_model(
828 &self,
829 issue_id: &str,
830 chunks: &[(usize, String, Vec<u8>)],
831 model_name: &str,
832 ) -> Result<()> {
833 let source_content_hash = self.compute_issue_embedding_content_hash(issue_id)?;
834 self.upsert_chunks_with_model_and_hash(issue_id, chunks, model_name, &source_content_hash)
835 }
836
837 pub fn upsert_chunks_with_model_and_hash(
838 &self,
839 issue_id: &str,
840 chunks: &[(usize, String, Vec<u8>)],
841 model_name: &str,
842 source_content_hash: &str,
843 ) -> Result<()> {
844 self.with_conn(|conn| {
845 let tx = conn.unchecked_transaction()?;
846 let (title, description): (String, Option<String>) = tx.query_row(
847 "SELECT title, description FROM issues WHERE id=?1",
848 rusqlite::params![issue_id],
849 |row| Ok((row.get(0)?, row.get(1)?)),
850 )?;
851 let current_hash =
852 crate::embedding::issue_content_hash(&title, description.as_deref());
853 if current_hash != source_content_hash {
854 anyhow::bail!(
855 "issue '{issue_id}' changed while its embeddings were being generated"
856 );
857 }
858 tx.execute(
859 "DELETE FROM chunks WHERE issue_id = ?1",
860 rusqlite::params![issue_id],
861 )?;
862 tx.execute(
863 "UPDATE issues SET embedding_content_hash=?2 WHERE id=?1",
864 rusqlite::params![issue_id, source_content_hash],
865 )?;
866 let mut stmt = tx.prepare(
867 "INSERT INTO chunks (issue_id, chunk_index, chunk_text, embedding, model_name, source_content_hash) VALUES (?1, ?2, ?3, ?4, ?5, ?6)"
868 )?;
869 for (idx, text, embedding) in chunks {
870 stmt.execute(rusqlite::params![
871 issue_id,
872 idx,
873 text,
874 embedding,
875 model_name,
876 source_content_hash
877 ])?;
878 }
879 drop(stmt);
880 tx.commit()?;
881 Ok(())
882 })
883 }
884
885 pub fn get_embedding_model(&self, issue_id: &str) -> Result<Option<String>> {
887 self.with_conn(|conn| {
888 let mut stmt =
889 conn.prepare("SELECT model_name FROM chunks WHERE issue_id = ?1 LIMIT 1")?;
890 let mut rows = stmt.query(rusqlite::params![issue_id])?;
891 if let Some(row) = rows.next()? {
892 let name: String = row.get(0)?;
893 Ok(if name.is_empty() { None } else { Some(name) })
894 } else {
895 Ok(None)
896 }
897 })
898 }
899
900 pub fn compute_issue_embedding_content_hash(&self, issue_id: &str) -> Result<String> {
901 self.with_conn(|conn| {
902 let (title, description): (String, Option<String>) = conn.query_row(
903 "SELECT title, description FROM issues WHERE id=?1",
904 rusqlite::params![issue_id],
905 |row| Ok((row.get(0)?, row.get(1)?)),
906 )?;
907 Ok(crate::embedding::issue_content_hash(
908 &title,
909 description.as_deref(),
910 ))
911 })
912 }
913
914 pub fn get_all_chunks(&self, workspace_id: &str) -> Result<Vec<Chunk>> {
915 self.with_conn(|conn| {
916 let mut stmt = conn.prepare(
917 "SELECT c.issue_id, c.embedding, i.identifier
918 FROM chunks c JOIN issues i ON c.issue_id = i.id
919 WHERE i.workspace_id = ?1",
920 )?;
921 let rows = stmt.query_map(rusqlite::params![workspace_id], |row| {
922 Ok(Chunk {
923 issue_id: row.get(0)?,
924 embedding: row.get(1)?,
925 identifier: row.get(2)?,
926 })
927 })?;
928 Ok(rows.collect::<std::result::Result<Vec<_>, _>>()?)
929 })
930 }
931
932 pub fn get_chunks_for_team(&self, team_key: &str, workspace_id: &str) -> Result<Vec<Chunk>> {
933 self.with_conn(|conn| {
934 let mut stmt = conn.prepare(
935 "SELECT c.issue_id, c.embedding, i.identifier
936 FROM chunks c JOIN issues i ON c.issue_id = i.id
937 WHERE i.team_key = ?1 AND i.workspace_id = ?2",
938 )?;
939 let rows = stmt.query_map(rusqlite::params![team_key, workspace_id], |row| {
940 Ok(Chunk {
941 issue_id: row.get(0)?,
942 embedding: row.get(1)?,
943 identifier: row.get(2)?,
944 })
945 })?;
946 Ok(rows.collect::<std::result::Result<Vec<_>, _>>()?)
947 })
948 }
949
950 pub fn count_embedded_issues(
951 &self,
952 team_key: Option<&str>,
953 workspace_id: &str,
954 ) -> Result<usize> {
955 self.with_conn(|conn| {
956 let count: usize = if let Some(team) = team_key {
957 conn.query_row(
958 "SELECT COUNT(DISTINCT c.issue_id) FROM chunks c JOIN issues i ON c.issue_id = i.id WHERE i.team_key = ?1 AND i.workspace_id = ?2",
959 rusqlite::params![team, workspace_id],
960 |row| row.get(0),
961 )?
962 } else {
963 conn.query_row(
964 "SELECT COUNT(DISTINCT c.issue_id) FROM chunks c JOIN issues i ON c.issue_id = i.id WHERE i.workspace_id = ?1",
965 rusqlite::params![workspace_id],
966 |row| row.get(0),
967 )?
968 };
969 Ok(count)
970 })
971 }
972
973 pub fn get_issues_needing_embedding(
974 &self,
975 team_key: Option<&str>,
976 force: bool,
977 workspace_id: &str,
978 ) -> Result<Vec<Issue>> {
979 self.get_issues_needing_embedding_for_model(team_key, force, workspace_id, None)
980 }
981
982 pub fn get_issues_needing_embedding_for_model(
983 &self,
984 team_key: Option<&str>,
985 force: bool,
986 workspace_id: &str,
987 model_name: Option<&str>,
988 ) -> Result<Vec<Issue>> {
989 self.with_conn(|conn| {
990 let mut stmt = conn.prepare(
991 "SELECT i.id, i.identifier, i.team_key, i.title, i.description,
992 i.state_name, i.state_type, i.priority, i.assignee_name,
993 i.project_name, i.labels_json, i.created_at, i.updated_at,
994 i.content_hash, i.synced_at, i.url, i.branch_name,
995 i.workspace_id, i.project_id, i.project_milestone_id,
996 i.project_milestone_name, i.cycle_id, i.cycle_name,
997 i.archived_at, i.embedding_content_hash,
998 c.model_name, c.source_content_hash, c.issue_id
999 FROM issues i
1000 JOIN issue_hydration_state h
1001 ON h.workspace_id=i.workspace_id AND h.issue_id=i.id
1002 AND h.resource='details' AND h.status='hydrated'
1003 LEFT JOIN (
1004 SELECT issue_id, MIN(model_name) AS model_name,
1005 MIN(source_content_hash) AS source_content_hash
1006 FROM chunks GROUP BY issue_id
1007 ) c ON c.issue_id=i.id
1008 WHERE i.workspace_id=?1 AND (?2 IS NULL OR i.team_key=?2)
1009 ORDER BY julianday(i.updated_at) DESC, i.identifier ASC",
1010 )?;
1011 let rows = stmt.query_map(rusqlite::params![workspace_id, team_key], |row| {
1012 Ok((
1013 Issue::from_row(row)?,
1014 row.get::<_, String>(24)?,
1015 row.get::<_, Option<String>>(25)?,
1016 row.get::<_, Option<String>>(26)?,
1017 row.get::<_, Option<String>>(27)?,
1018 ))
1019 })?;
1020 let mut issues = Vec::new();
1021 for row in rows {
1022 let (issue, current_hash, embedded_model, embedded_hash, chunk_issue_id) = row?;
1023 let needs_embedding = force
1024 || chunk_issue_id.is_none()
1025 || current_hash.is_empty()
1026 || embedded_hash.as_deref().unwrap_or_default().is_empty()
1027 || embedded_hash.as_deref() != Some(current_hash.as_str())
1028 || model_name.is_some_and(|model| embedded_model.as_deref() != Some(model));
1029 if needs_embedding {
1030 issues.push(issue);
1031 }
1032 }
1033 Ok(issues)
1034 })
1035 }
1036
1037 pub fn get_comments(&self, issue_id: &str) -> Result<Vec<Comment>> {
1040 self.with_conn(|conn| {
1041 let mut stmt = conn.prepare(
1042 "SELECT id, issue_id, body, user_name, created_at, updated_at, parent_id, url, workspace_id
1043 FROM comments
1044 WHERE issue_id = ?1
1045 ORDER BY created_at"
1046 )?;
1047 let rows = stmt.query_map(rusqlite::params![issue_id], |row| {
1048 Ok(Comment {
1049 id: row.get(0)?,
1050 issue_id: row.get(1)?,
1051 body: row.get(2)?,
1052 user_name: row.get(3)?,
1053 created_at: row.get(4)?,
1054 updated_at: row.get(5)?,
1055 parent_id: row.get(6)?,
1056 url: row.get(7)?,
1057 workspace_id: row.get(8)?,
1058 })
1059 })?;
1060 Ok(rows.collect::<std::result::Result<Vec<_>, _>>()?)
1061 })
1062 }
1063
1064 pub fn replace_issue_comments(
1065 &self,
1066 issue_id: &str,
1067 workspace_id: &str,
1068 comments: &[Comment],
1069 ) -> Result<()> {
1070 self.with_conn(|conn| {
1071 let tx = conn.unchecked_transaction()?;
1072 tx.execute(
1073 "DELETE FROM comments WHERE issue_id = ?1 AND workspace_id = ?2",
1074 rusqlite::params![issue_id, workspace_id],
1075 )?;
1076 for comment in comments {
1077 tx.execute(
1078 "INSERT INTO comments
1079 (id, issue_id, body, user_name, created_at, workspace_id, updated_at, parent_id, url)
1080 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)
1081 ON CONFLICT(id) DO UPDATE SET
1082 issue_id=excluded.issue_id,
1083 body=excluded.body,
1084 user_name=excluded.user_name,
1085 created_at=excluded.created_at,
1086 workspace_id=excluded.workspace_id,
1087 updated_at=excluded.updated_at,
1088 parent_id=excluded.parent_id,
1089 url=excluded.url",
1090 rusqlite::params![
1091 comment.id,
1092 comment.issue_id,
1093 comment.body,
1094 comment.user_name,
1095 comment.created_at,
1096 workspace_id,
1097 comment.updated_at,
1098 comment.parent_id,
1099 comment.url,
1100 ],
1101 )?;
1102 }
1103 tx.commit()?;
1104 Ok(())
1105 })
1106 }
1107
1108 pub fn get_comment_sync_state(&self, issue_id: &str) -> Result<CommentSyncState> {
1109 self.with_conn(|conn| {
1110 let mut stmt = conn.prepare(
1111 "SELECT status, sync_error, synced_at
1112 FROM comment_sync_state
1113 WHERE issue_id = ?1",
1114 )?;
1115 let mut rows = stmt.query(rusqlite::params![issue_id])?;
1116 if let Some(row) = rows.next()? {
1117 Ok(CommentSyncState {
1118 status: row.get(0)?,
1119 sync_error: row.get(1)?,
1120 synced_at: row.get(2)?,
1121 })
1122 } else {
1123 Ok(CommentSyncState::not_synced())
1124 }
1125 })
1126 }
1127
1128 pub fn mark_comments_synced(
1129 &self,
1130 issue_id: &str,
1131 workspace_id: &str,
1132 comment_count: usize,
1133 ) -> Result<()> {
1134 let status = if comment_count == 0 {
1135 "none_found"
1136 } else {
1137 "synced"
1138 };
1139 self.set_comment_sync_state(issue_id, workspace_id, status, None)
1140 }
1141
1142 pub fn mark_comments_sync_failed(
1143 &self,
1144 issue_id: &str,
1145 workspace_id: &str,
1146 status: &str,
1147 error: &str,
1148 ) -> Result<()> {
1149 self.set_comment_sync_state(issue_id, workspace_id, status, Some(error))
1150 }
1151
1152 fn set_comment_sync_state(
1153 &self,
1154 issue_id: &str,
1155 workspace_id: &str,
1156 status: &str,
1157 error: Option<&str>,
1158 ) -> Result<()> {
1159 self.with_conn(|conn| {
1160 conn.execute(
1161 "INSERT INTO comment_sync_state (issue_id, workspace_id, status, sync_error, synced_at)
1162 VALUES (?1, ?2, ?3, ?4, datetime('now'))
1163 ON CONFLICT(issue_id) DO UPDATE SET
1164 workspace_id=excluded.workspace_id,
1165 status=excluded.status,
1166 sync_error=excluded.sync_error,
1167 synced_at=datetime('now')",
1168 rusqlite::params![issue_id, workspace_id, status, error],
1169 )?;
1170 Ok(())
1171 })
1172 }
1173
1174 pub fn get_sync_cursor(&self, workspace_id: &str, team_key: &str) -> Result<Option<String>> {
1177 self.get_synced_through_at(workspace_id, team_key)
1178 }
1179
1180 pub fn get_synced_through_at(
1181 &self,
1182 workspace_id: &str,
1183 team_key: &str,
1184 ) -> Result<Option<String>> {
1185 self.with_conn(|conn| {
1186 let mut stmt = conn.prepare(
1187 "SELECT COALESCE(synced_through_at, last_updated_at)
1188 FROM sync_state WHERE workspace_id = ?1 AND team_key = ?2",
1189 )?;
1190 let mut rows = stmt.query(rusqlite::params![workspace_id, team_key])?;
1191 if let Some(row) = rows.next()? {
1192 Ok(Some(row.get(0)?))
1193 } else {
1194 Ok(None)
1195 }
1196 })
1197 }
1198
1199 pub fn set_sync_cursor(
1200 &self,
1201 workspace_id: &str,
1202 team_key: &str,
1203 last_updated_at: &str,
1204 ) -> Result<()> {
1205 self.with_conn(|conn| {
1206 conn.execute(
1207 "INSERT INTO sync_state (
1208 workspace_id, team_key, last_updated_at, synced_through_at,
1209 full_sync_done, last_synced_at
1210 ) VALUES (?1, ?2, ?3, ?3, 1, datetime('now'))
1211 ON CONFLICT(workspace_id, team_key) DO UPDATE SET
1212 last_updated_at=excluded.last_updated_at,
1213 synced_through_at=excluded.synced_through_at,
1214 full_sync_done=1,
1215 last_synced_at=datetime('now')",
1216 rusqlite::params![workspace_id, team_key, last_updated_at],
1217 )?;
1218 Ok(())
1219 })
1220 }
1221
1222 pub fn is_full_sync_done(&self, workspace_id: &str, team_key: &str) -> Result<bool> {
1223 self.with_conn(|conn| {
1224 let mut stmt = conn.prepare(
1225 "SELECT full_sync_done FROM sync_state WHERE workspace_id = ?1 AND team_key = ?2",
1226 )?;
1227 let mut rows = stmt.query(rusqlite::params![workspace_id, team_key])?;
1228 if let Some(row) = rows.next()? {
1229 let done: bool = row.get(0)?;
1230 Ok(done)
1231 } else {
1232 Ok(false)
1233 }
1234 })
1235 }
1236
1237 pub fn get_last_synced_at(&self, workspace_id: &str, team_key: &str) -> Result<Option<String>> {
1239 self.with_conn(|conn| {
1240 let mut stmt = conn.prepare(
1241 "SELECT last_synced_at FROM sync_state WHERE workspace_id = ?1 AND team_key = ?2",
1242 )?;
1243 let mut rows = stmt.query(rusqlite::params![workspace_id, team_key])?;
1244 if let Some(row) = rows.next()? {
1245 Ok(row.get(0)?)
1246 } else {
1247 Ok(None)
1248 }
1249 })
1250 }
1251
1252 pub fn get_metadata(&self, key: &str) -> Result<Option<String>> {
1255 self.with_conn(|conn| {
1256 let mut stmt = conn.prepare("SELECT value FROM metadata WHERE key = ?1")?;
1257 let mut rows = stmt.query(rusqlite::params![key])?;
1258 if let Some(row) = rows.next()? {
1259 Ok(Some(row.get(0)?))
1260 } else {
1261 Ok(None)
1262 }
1263 })
1264 }
1265
1266 pub fn set_metadata(&self, key: &str, value: &str) -> Result<()> {
1267 self.with_conn(|conn| {
1268 conn.execute(
1269 "INSERT INTO metadata (key, value) VALUES (?1, ?2) ON CONFLICT(key) DO UPDATE SET value=excluded.value",
1270 rusqlite::params![key, value],
1271 )?;
1272 Ok(())
1273 })
1274 }
1275
1276 pub fn list_synced_teams(&self, workspace_id: &str) -> Result<Vec<TeamSummary>> {
1281 self.with_conn(|conn| {
1282 let mut stmt = conn.prepare(
1283 "SELECT i.team_key,
1284 COUNT(DISTINCT i.id) AS issue_count,
1285 COUNT(DISTINCT c.issue_id) AS embedded_count,
1286 s.last_synced_at
1287 FROM issues i
1288 LEFT JOIN chunks c ON i.id = c.issue_id
1289 LEFT JOIN sync_state s ON i.team_key = s.team_key AND s.workspace_id = ?1
1290 WHERE i.workspace_id = ?1
1291 GROUP BY i.team_key
1292 ORDER BY i.team_key",
1293 )?;
1294 let rows = stmt.query_map(rusqlite::params![workspace_id], |row| {
1295 Ok(TeamSummary {
1296 key: row.get(0)?,
1297 issue_count: row.get(1)?,
1298 embedded_count: row.get(2)?,
1299 last_synced_at: row.get(3)?,
1300 })
1301 })?;
1302 Ok(rows.collect::<std::result::Result<Vec<_>, _>>()?)
1303 })
1304 }
1305
1306 pub fn fts_search(
1307 &self,
1308 query: &str,
1309 limit: usize,
1310 workspace_id: &str,
1311 ) -> Result<Vec<FtsResult>> {
1312 self.fts_search_filtered(query, limit, workspace_id, None)
1313 }
1314
1315 pub fn fts_search_filtered(
1316 &self,
1317 query: &str,
1318 limit: usize,
1319 workspace_id: &str,
1320 label_ids: Option<&[String]>,
1321 ) -> Result<Vec<FtsResult>> {
1322 self.with_conn(|conn| {
1323 let mut params: Vec<Box<dyn rusqlite::types::ToSql>> = vec![
1324 Box::new(query.to_string()),
1325 Box::new(limit as i64),
1326 Box::new(workspace_id.to_string()),
1327 ];
1328
1329 let label_clause = if let Some(ids) = label_ids.filter(|ids| !ids.is_empty()) {
1330 let (frag, mut lp) = Self::label_filter_fragment(ids, params.len() + 1, "i");
1331 params.append(&mut lp);
1332 format!(" AND {frag}")
1333 } else {
1334 String::new()
1335 };
1336
1337 let sql = format!(
1338 "SELECT i.id, i.identifier, i.title, i.state_name, i.priority, bm25(issues_fts) as rank
1339 FROM issues_fts f
1340 JOIN issues i ON f.rowid = i.rowid
1341 WHERE issues_fts MATCH ?1 AND i.workspace_id = ?3{label_clause}
1342 ORDER BY rank
1343 LIMIT ?2"
1344 );
1345 let mut stmt = conn.prepare(&sql)?;
1346 let param_refs: Vec<&dyn rusqlite::types::ToSql> = params.iter().map(|p| p.as_ref()).collect();
1347 let rows = stmt.query_map(param_refs.as_slice(), |row| {
1348 Ok(FtsResult {
1349 issue_id: row.get(0)?,
1350 identifier: row.get(1)?,
1351 title: row.get(2)?,
1352 state_name: row.get(3)?,
1353 priority: row.get(4)?,
1354 bm25_score: row.get(5)?,
1355 })
1356 })?;
1357 Ok(rows.collect::<std::result::Result<Vec<_>, _>>()?)
1358 })
1359 }
1360}
1361
1362fn default_workspace_id() -> String {
1365 "default".to_string()
1366}
1367
1368#[derive(Debug, Clone, Serialize, Deserialize)]
1369pub struct Issue {
1370 pub id: String,
1371 pub identifier: String,
1372 pub team_key: String,
1373 pub title: String,
1374 pub description: Option<String>,
1375 pub state_name: String,
1376 pub state_type: String,
1377 pub priority: i32,
1378 pub assignee_name: Option<String>,
1379 pub project_name: Option<String>,
1380 pub labels_json: String,
1381 pub created_at: String,
1382 pub updated_at: String,
1383 pub content_hash: String,
1384 pub synced_at: Option<String>,
1385 pub url: String,
1386 pub branch_name: Option<String>,
1387 #[serde(default = "default_workspace_id")]
1388 pub workspace_id: String,
1389 pub project_id: Option<String>,
1390 pub project_milestone_id: Option<String>,
1391 pub project_milestone_name: Option<String>,
1392 pub cycle_id: Option<String>,
1393 pub cycle_name: Option<String>,
1394 pub archived_at: Option<String>,
1395}
1396
1397impl Issue {
1398 pub fn from_row(row: &rusqlite::Row) -> rusqlite::Result<Self> {
1399 Ok(Self {
1400 id: row.get(0)?,
1401 identifier: row.get(1)?,
1402 team_key: row.get(2)?,
1403 title: row.get(3)?,
1404 description: row.get(4)?,
1405 state_name: row.get(5)?,
1406 state_type: row.get(6)?,
1407 priority: row.get(7)?,
1408 assignee_name: row.get(8)?,
1409 project_name: row.get(9)?,
1410 labels_json: row.get(10)?,
1411 created_at: row.get(11)?,
1412 updated_at: row.get(12)?,
1413 content_hash: row.get(13)?,
1414 synced_at: row.get(14)?,
1415 url: row.get(15)?,
1416 branch_name: row.get(16).unwrap_or(None),
1417 workspace_id: row.get(17).unwrap_or_else(|_| "default".to_string()),
1418 project_id: row.get(18).unwrap_or(None),
1419 project_milestone_id: row.get(19).unwrap_or(None),
1420 project_milestone_name: row.get(20).unwrap_or(None),
1421 cycle_id: row.get(21).unwrap_or(None),
1422 cycle_name: row.get(22).unwrap_or(None),
1423 archived_at: row.get(23).unwrap_or(None),
1424 })
1425 }
1426
1427 pub fn labels(&self) -> Vec<String> {
1428 serde_json::from_str(&self.labels_json).unwrap_or_default()
1429 }
1430
1431 pub fn priority_label(&self) -> &str {
1432 match self.priority {
1433 0 => "No priority",
1434 1 => "Urgent",
1435 2 => "High",
1436 3 => "Medium",
1437 4 => "Low",
1438 _ => "Unknown",
1439 }
1440 }
1441}
1442
1443#[derive(Debug, Clone, Serialize, Deserialize)]
1444pub struct Relation {
1445 pub id: String,
1446 pub issue_id: String,
1447 pub related_issue_id: String,
1448 pub related_issue_identifier: String,
1449 pub relation_type: String,
1450}
1451
1452#[derive(Debug, Clone, Serialize, Deserialize)]
1453pub struct EnrichedRelation {
1454 pub relation_id: String,
1455 pub relation_type: String,
1456 pub issue_identifier: String,
1457 pub issue_title: String,
1458 pub issue_state: String,
1459 pub issue_url: String,
1460}
1461
1462#[derive(Debug, Clone)]
1463pub struct Chunk {
1464 pub issue_id: String,
1465 pub embedding: Vec<u8>,
1466 pub identifier: String,
1467}
1468
1469#[derive(Debug, Clone, Serialize, Deserialize)]
1470pub struct Comment {
1471 pub id: String,
1472 pub issue_id: String,
1473 pub body: String,
1474 pub user_name: Option<String>,
1475 pub created_at: String,
1476 pub updated_at: Option<String>,
1477 pub parent_id: Option<String>,
1478 pub url: Option<String>,
1479 #[serde(default = "default_workspace_id")]
1480 pub workspace_id: String,
1481}
1482
1483#[derive(Debug, Clone, Serialize, Deserialize)]
1484pub struct CommentSyncState {
1485 pub status: String,
1486 pub sync_error: Option<String>,
1487 pub synced_at: Option<String>,
1488}
1489
1490impl CommentSyncState {
1491 pub fn not_synced() -> Self {
1492 Self {
1493 status: "not_synced".to_string(),
1494 sync_error: None,
1495 synced_at: None,
1496 }
1497 }
1498}
1499
1500#[derive(Debug, Clone)]
1501pub struct FtsResult {
1502 pub issue_id: String,
1503 pub identifier: String,
1504 pub title: String,
1505 pub state_name: String,
1506 pub priority: i32,
1507 pub bm25_score: f64,
1508}
1509
1510#[derive(Debug, Clone)]
1511pub struct IssueSummary {
1512 pub id: String,
1513 pub identifier: String,
1514 pub team_key: String,
1515 pub title: String,
1516 pub state_name: String,
1517 pub state_type: String,
1518 pub priority: i32,
1519 pub project_name: Option<String>,
1520 pub labels: Vec<String>,
1521 pub updated_at: String,
1522 pub url: String,
1523 pub has_description: bool,
1524 pub has_embedding: bool,
1525}
1526
1527#[derive(Debug, Clone)]
1528pub struct TeamSummary {
1529 pub key: String,
1530 pub issue_count: usize,
1531 pub embedded_count: usize,
1532 pub last_synced_at: Option<String>,
1533}
1534
1535#[derive(Debug, Clone, Serialize, Deserialize)]
1536pub struct WorkspaceRow {
1537 pub id: String,
1538 pub linear_org_id: Option<String>,
1539 pub display_name: Option<String>,
1540 pub created_at: String,
1541}
1542
1543#[derive(Debug, Clone, Serialize, Deserialize)]
1544pub struct Label {
1545 pub id: String,
1546 pub workspace_id: String,
1547 pub name: String,
1548 pub color: Option<String>,
1549 pub parent_id: Option<String>,
1550}
1551
1552#[cfg(test)]
1553mod tests {
1554 use super::test_helpers::*;
1555 use super::Comment;
1556
1557 #[test]
1558 fn count_embedded_issues_empty_db() {
1559 let (db, _dir) = test_db();
1560 assert_eq!(db.count_embedded_issues(None, "default").unwrap(), 0);
1561 }
1562
1563 #[test]
1564 fn comment_sync_state_defaults_to_not_synced() {
1565 let (db, _dir) = test_db();
1566 let issue = make_issue("TST-1", "TST");
1567 db.upsert_issue(&issue).unwrap();
1568
1569 let state = db.get_comment_sync_state(&issue.id).unwrap();
1570 assert_eq!(state.status, "not_synced");
1571 assert!(state.synced_at.is_none());
1572 assert!(state.sync_error.is_none());
1573 }
1574
1575 #[test]
1576 fn replace_issue_comments_preserves_thread_metadata() {
1577 let (db, _dir) = test_db();
1578 let issue = make_issue("TST-1", "TST");
1579 db.upsert_issue(&issue).unwrap();
1580
1581 db.replace_issue_comments(
1582 &issue.id,
1583 "default",
1584 &[Comment {
1585 id: "comment-1".to_string(),
1586 issue_id: issue.id.clone(),
1587 body: "fixed in linked PR".to_string(),
1588 user_name: Some("Ada".to_string()),
1589 created_at: "2026-01-03T00:00:00Z".to_string(),
1590 updated_at: Some("2026-01-03T01:00:00Z".to_string()),
1591 parent_id: Some("parent-1".to_string()),
1592 url: Some("https://linear.app/comment/comment-1".to_string()),
1593 workspace_id: "default".to_string(),
1594 }],
1595 )
1596 .unwrap();
1597 db.mark_comments_synced(&issue.id, "default", 1).unwrap();
1598
1599 let comments = db.get_comments(&issue.id).unwrap();
1600 assert_eq!(comments.len(), 1);
1601 assert_eq!(comments[0].parent_id.as_deref(), Some("parent-1"));
1602 assert_eq!(
1603 comments[0].updated_at.as_deref(),
1604 Some("2026-01-03T01:00:00Z")
1605 );
1606 assert_eq!(
1607 comments[0].url.as_deref(),
1608 Some("https://linear.app/comment/comment-1")
1609 );
1610 assert_eq!(
1611 db.get_comment_sync_state(&issue.id).unwrap().status,
1612 "synced"
1613 );
1614 }
1615
1616 #[test]
1617 fn empty_comment_sync_records_none_found() {
1618 let (db, _dir) = test_db();
1619 let issue = make_issue("TST-1", "TST");
1620 db.upsert_issue(&issue).unwrap();
1621
1622 db.replace_issue_comments(&issue.id, "default", &[])
1623 .unwrap();
1624 db.mark_comments_synced(&issue.id, "default", 0).unwrap();
1625
1626 assert!(db.get_comments(&issue.id).unwrap().is_empty());
1627 let state = db.get_comment_sync_state(&issue.id).unwrap();
1628 assert_eq!(state.status, "none_found");
1629 assert!(state.synced_at.is_some());
1630 }
1631
1632 #[test]
1633 fn count_embedded_issues_with_data() {
1634 let (db, _dir) = test_db();
1635
1636 let issue1 = make_issue("TST-1", "TST");
1637 let issue2 = make_issue("TST-2", "TST");
1638 let issue3 = make_issue("OTH-1", "OTH");
1639 db.upsert_issue(&issue1).unwrap();
1640 db.upsert_issue(&issue2).unwrap();
1641 db.upsert_issue(&issue3).unwrap();
1642
1643 db.upsert_chunks(&issue1.id, &[(0, "chunk".into(), fake_embedding(768))])
1645 .unwrap();
1646 db.upsert_chunks(&issue3.id, &[(0, "chunk".into(), fake_embedding(768))])
1647 .unwrap();
1648
1649 assert_eq!(db.count_embedded_issues(None, "default").unwrap(), 2);
1651 assert_eq!(db.count_embedded_issues(Some("TST"), "default").unwrap(), 1);
1653 assert_eq!(db.count_embedded_issues(Some("OTH"), "default").unwrap(), 1);
1654 assert_eq!(
1655 db.count_embedded_issues(Some("NONE"), "default").unwrap(),
1656 0
1657 );
1658 }
1659
1660 #[test]
1661 fn get_field_completeness_empty_db() {
1662 let (db, _dir) = test_db();
1663 let (total, desc, pri, labels, proj) = db.get_field_completeness(None, "default").unwrap();
1664 assert_eq!(total, 0);
1665 assert_eq!(desc, 0);
1666 assert_eq!(pri, 0);
1667 assert_eq!(labels, 0);
1668 assert_eq!(proj, 0);
1669 }
1670
1671 #[test]
1672 fn get_field_completeness_with_data() {
1673 let (db, _dir) = test_db();
1674
1675 let mut full = make_issue("TST-1", "TST");
1677 full.description = Some("Has desc".into());
1678 full.priority = 2;
1679 full.labels_json = r#"["bug"]"#.into();
1680 full.project_name = Some("Proj".into());
1681 db.upsert_issue(&full).unwrap();
1682
1683 let mut sparse = make_issue("TST-2", "TST");
1685 sparse.description = None;
1686 sparse.priority = 0;
1687 sparse.labels_json = "[]".into();
1688 sparse.project_name = None;
1689 db.upsert_issue(&sparse).unwrap();
1690
1691 let mut other = make_issue("OTH-1", "OTH");
1693 other.description = Some("Has desc".into());
1694 other.priority = 0;
1695 other.labels_json = "[]".into();
1696 other.project_name = None;
1697 db.upsert_issue(&other).unwrap();
1698
1699 let (total, desc, pri, labels, proj) = db.get_field_completeness(None, "default").unwrap();
1701 assert_eq!(total, 3);
1702 assert_eq!(desc, 2); assert_eq!(pri, 1); assert_eq!(labels, 1); assert_eq!(proj, 1); let (total, desc, pri, labels, proj) =
1709 db.get_field_completeness(Some("TST"), "default").unwrap();
1710 assert_eq!(total, 2);
1711 assert_eq!(desc, 1);
1712 assert_eq!(pri, 1);
1713 assert_eq!(labels, 1);
1714 assert_eq!(proj, 1);
1715 }
1716
1717 #[test]
1718 fn list_all_issues_pagination_and_filter() {
1719 let (db, _dir) = test_db();
1720
1721 for i in 1..=5 {
1722 let mut issue = make_issue(&format!("TST-{i}"), "TST");
1723 issue.updated_at = format!("2026-01-0{i}T00:00:00Z");
1724 db.upsert_issue(&issue).unwrap();
1725 }
1726 let mut other = make_issue("OTH-1", "OTH");
1727 other.updated_at = "2026-01-06T00:00:00Z".to_string();
1728 db.upsert_issue(&other).unwrap();
1729
1730 let page1 = db.list_all_issues(None, None, 3, 0, "default").unwrap();
1732 assert_eq!(page1.len(), 3);
1733 assert_eq!(page1[0].identifier, "OTH-1");
1735
1736 let page2 = db.list_all_issues(None, None, 3, 3, "default").unwrap();
1738 assert_eq!(page2.len(), 3);
1739
1740 let page3 = db.list_all_issues(None, None, 3, 6, "default").unwrap();
1742 assert_eq!(page3.len(), 0);
1743
1744 let tst = db
1746 .list_all_issues(Some("TST"), None, 10, 0, "default")
1747 .unwrap();
1748 assert_eq!(tst.len(), 5);
1749
1750 let filtered = db
1752 .list_all_issues(None, Some("TST-3"), 10, 0, "default")
1753 .unwrap();
1754 assert_eq!(filtered.len(), 1);
1755 assert_eq!(filtered[0].identifier, "TST-3");
1756
1757 let title_match = db
1759 .list_all_issues(None, Some("Test issue OTH"), 10, 0, "default")
1760 .unwrap();
1761 assert_eq!(title_match.len(), 1);
1762 }
1763
1764 #[test]
1765 fn list_all_issues_has_embedding_flag() {
1766 let (db, _dir) = test_db();
1767
1768 let issue1 = make_issue("TST-1", "TST");
1769 let issue2 = make_issue("TST-2", "TST");
1770 db.upsert_issue(&issue1).unwrap();
1771 db.upsert_issue(&issue2).unwrap();
1772
1773 db.upsert_chunks(&issue1.id, &[(0, "chunk".into(), fake_embedding(768))])
1775 .unwrap();
1776
1777 let issues = db.list_all_issues(None, None, 10, 0, "default").unwrap();
1778 let by_id: std::collections::HashMap<_, _> =
1779 issues.iter().map(|i| (i.identifier.as_str(), i)).collect();
1780
1781 assert!(by_id["TST-1"].has_embedding);
1782 assert!(!by_id["TST-2"].has_embedding);
1783 }
1784
1785 #[test]
1786 fn list_synced_teams_empty_db() {
1787 let (db, _dir) = test_db();
1788 let teams = db.list_synced_teams("default").unwrap();
1789 assert!(teams.is_empty());
1790 }
1791
1792 #[test]
1793 fn list_synced_teams_with_data() {
1794 let (db, _dir) = test_db();
1795
1796 for i in 1..=3 {
1798 let issue = make_issue(&format!("TST-{i}"), "TST");
1799 db.upsert_issue(&issue).unwrap();
1800 if i <= 2 {
1801 db.upsert_chunks(&issue.id, &[(0, "chunk".into(), fake_embedding(768))])
1803 .unwrap();
1804 }
1805 }
1806 let other = make_issue("OTH-1", "OTH");
1807 db.upsert_issue(&other).unwrap();
1808
1809 let teams = db.list_synced_teams("default").unwrap();
1810 assert_eq!(teams.len(), 2);
1811
1812 let by_key: std::collections::HashMap<_, _> =
1814 teams.iter().map(|t| (t.key.as_str(), t)).collect();
1815
1816 assert_eq!(by_key["TST"].issue_count, 3);
1817 assert_eq!(by_key["TST"].embedded_count, 2);
1818 assert_eq!(by_key["OTH"].issue_count, 1);
1819 assert_eq!(by_key["OTH"].embedded_count, 0);
1820 }
1821
1822 #[test]
1823 fn list_synced_teams_includes_last_synced_at() {
1824 let (db, _dir) = test_db();
1825
1826 let issue = make_issue("TST-1", "TST");
1827 db.upsert_issue(&issue).unwrap();
1828
1829 let teams = db.list_synced_teams("default").unwrap();
1831 assert_eq!(teams.len(), 1);
1832 assert!(teams[0].last_synced_at.is_none());
1833
1834 db.set_sync_cursor("default", "TST", "2026-01-01T00:00:00Z")
1836 .unwrap();
1837 let teams = db.list_synced_teams("default").unwrap();
1838 assert!(teams[0].last_synced_at.is_some());
1839 }
1840
1841 #[test]
1842 fn list_synced_teams_multi_chunk_issue() {
1843 let (db, _dir) = test_db();
1844
1845 let issue = make_issue("TST-1", "TST");
1846 db.upsert_issue(&issue).unwrap();
1847 db.upsert_chunks(
1849 &issue.id,
1850 &[
1851 (0, "chunk0".into(), fake_embedding(768)),
1852 (1, "chunk1".into(), fake_embedding(768)),
1853 (2, "chunk2".into(), fake_embedding(768)),
1854 ],
1855 )
1856 .unwrap();
1857
1858 let teams = db.list_synced_teams("default").unwrap();
1859 assert_eq!(teams.len(), 1);
1860 assert_eq!(teams[0].issue_count, 1); assert_eq!(teams[0].embedded_count, 1);
1862 }
1863
1864 #[test]
1865 fn workspace_crud() {
1866 let (db, _dir) = test_db();
1867
1868 let ws = db.get_workspace("default").unwrap();
1870 assert!(ws.is_some());
1871
1872 db.upsert_workspace("work", None, None).unwrap();
1874 let ws = db.get_workspace("work").unwrap().unwrap();
1875 assert_eq!(ws.id, "work");
1876 assert!(ws.linear_org_id.is_none());
1877
1878 db.upsert_workspace("work", Some("org-123"), Some("Work Org"))
1880 .unwrap();
1881 let ws = db.get_workspace("work").unwrap().unwrap();
1882 assert_eq!(ws.linear_org_id.as_deref(), Some("org-123"));
1883 assert_eq!(ws.display_name.as_deref(), Some("Work Org"));
1884
1885 let all = db.list_workspaces().unwrap();
1887 assert_eq!(all.len(), 2);
1888
1889 db.delete_workspace("work").unwrap();
1891 let ws = db.get_workspace("work").unwrap();
1892 assert!(ws.is_none());
1893 }
1894
1895 #[test]
1896 fn issues_isolated_by_workspace() {
1897 let (db, _dir) = test_db();
1898
1899 db.upsert_workspace("work", None, None).unwrap();
1901
1902 let mut issue1 = make_issue("TST-1", "TST");
1904 issue1.workspace_id = "default".to_string();
1905 issue1.priority = 0;
1906 db.upsert_issue(&issue1).unwrap();
1907
1908 let mut issue2 = make_issue("TST-2", "TST");
1910 issue2.id = "id-2".to_string();
1911 issue2.workspace_id = "work".to_string();
1912 issue2.priority = 0;
1913 db.upsert_issue(&issue2).unwrap();
1914
1915 assert_eq!(db.count_issues(None, "default").unwrap(), 1);
1917 assert_eq!(db.count_issues(None, "work").unwrap(), 1);
1918
1919 let default_unpri = db.get_unprioritized_issues(None, false, "default").unwrap();
1921 assert_eq!(default_unpri.len(), 1);
1922 assert_eq!(default_unpri[0].identifier, "TST-1");
1923
1924 let work_unpri = db.get_unprioritized_issues(None, false, "work").unwrap();
1925 assert_eq!(work_unpri.len(), 1);
1926 assert_eq!(work_unpri[0].identifier, "TST-2");
1927 }
1928
1929 #[test]
1930 fn sync_state_isolated_by_workspace() {
1931 let (db, _dir) = test_db();
1932 db.upsert_workspace("work", None, None).unwrap();
1933
1934 db.set_sync_cursor("default", "TST", "2024-01-01T00:00:00Z")
1936 .unwrap();
1937 db.set_sync_cursor("work", "TST", "2024-06-01T00:00:00Z")
1938 .unwrap();
1939
1940 assert_eq!(
1941 db.get_sync_cursor("default", "TST").unwrap().as_deref(),
1942 Some("2024-01-01T00:00:00Z")
1943 );
1944 assert_eq!(
1945 db.get_sync_cursor("work", "TST").unwrap().as_deref(),
1946 Some("2024-06-01T00:00:00Z")
1947 );
1948
1949 assert!(db.is_full_sync_done("default", "TST").unwrap());
1950 assert!(db.is_full_sync_done("work", "TST").unwrap());
1951 assert!(!db.is_full_sync_done("default", "OTHER").unwrap());
1952 }
1953
1954 #[test]
1955 fn list_synced_teams_workspace_scoped() {
1956 let (db, _dir) = test_db();
1957 db.upsert_workspace("work", None, None).unwrap();
1958
1959 let mut issue1 = make_issue("TST-1", "TST");
1960 issue1.workspace_id = "default".to_string();
1961 db.upsert_issue(&issue1).unwrap();
1962
1963 let mut issue2 = make_issue("WRK-1", "WRK");
1964 issue2.id = "id-wrk".to_string();
1965 issue2.workspace_id = "work".to_string();
1966 db.upsert_issue(&issue2).unwrap();
1967
1968 let default_teams = db.list_synced_teams("default").unwrap();
1969 assert_eq!(default_teams.len(), 1);
1970 assert_eq!(default_teams[0].key, "TST");
1971
1972 let work_teams = db.list_synced_teams("work").unwrap();
1973 assert_eq!(work_teams.len(), 1);
1974 assert_eq!(work_teams[0].key, "WRK");
1975 }
1976
1977 #[test]
1978 fn migration_8_creates_label_tables_and_resets_sync_state() {
1979 let dir = tempfile::tempdir().unwrap();
1980 let path = dir.path().join("test.db");
1981 let conn = rusqlite::Connection::open(&path).unwrap();
1982
1983 conn.execute_batch("PRAGMA journal_mode=WAL; PRAGMA foreign_keys=ON;")
1984 .unwrap();
1985
1986 crate::db::schema::run_migrations(&conn).unwrap();
1988
1989 conn.execute("DELETE FROM schema_version WHERE version >= 8", [])
1991 .unwrap();
1992
1993 conn.execute(
1995 "INSERT INTO sync_state (workspace_id, team_key, last_updated_at, full_sync_done, last_synced_at)
1996 VALUES ('default', 'ENG', '2026-04-01T00:00:00Z', 1, '2026-04-01T00:00:00Z')",
1997 [],
1998 )
1999 .unwrap();
2000
2001 let full_done_before: i64 = conn
2003 .query_row(
2004 "SELECT full_sync_done FROM sync_state WHERE workspace_id='default' AND team_key='ENG'",
2005 [],
2006 |r| r.get(0),
2007 )
2008 .unwrap();
2009 assert_eq!(full_done_before, 1);
2010
2011 crate::db::schema::run_migrations(&conn).unwrap();
2013
2014 let labels_count: i64 = conn
2016 .query_row(
2017 "SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name='labels'",
2018 [],
2019 |r| r.get(0),
2020 )
2021 .unwrap();
2022 assert_eq!(labels_count, 1);
2023 let join_count: i64 = conn
2024 .query_row(
2025 "SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name='issue_labels'",
2026 [],
2027 |r| r.get(0),
2028 )
2029 .unwrap();
2030 assert_eq!(join_count, 1);
2031
2032 let full_done: i64 = conn
2034 .query_row(
2035 "SELECT full_sync_done FROM sync_state WHERE workspace_id='default' AND team_key='ENG'",
2036 [],
2037 |r| r.get(0),
2038 )
2039 .unwrap();
2040 assert_eq!(full_done, 0);
2041 let last_updated: String = conn
2042 .query_row(
2043 "SELECT last_updated_at FROM sync_state WHERE workspace_id='default' AND team_key='ENG'",
2044 [],
2045 |r| r.get(0),
2046 )
2047 .unwrap();
2048 assert_eq!(last_updated, "1970-01-01T00:00:00Z");
2049 }
2050
2051 #[test]
2052 fn migration_10_repairs_missing_project_join_tables() {
2053 let dir = tempfile::tempdir().unwrap();
2054 let path = dir.path().join("test.db");
2055 let conn = rusqlite::Connection::open(&path).unwrap();
2056
2057 crate::db::schema::run_migrations(&conn).unwrap();
2058 conn.execute("DROP TABLE project_labels", []).unwrap();
2059
2060 let version: i64 = conn
2061 .query_row("SELECT MAX(version) FROM schema_version", [], |row| {
2062 row.get(0)
2063 })
2064 .unwrap();
2065 assert_eq!(version, 13);
2066
2067 crate::db::schema::run_migrations(&conn).unwrap();
2068
2069 let table_count: i64 = conn
2070 .query_row(
2071 "SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name='project_labels'",
2072 [],
2073 |row| row.get(0),
2074 )
2075 .unwrap();
2076 assert_eq!(table_count, 1);
2077 }
2078
2079 #[test]
2080 fn upsert_label_inserts_and_renames_in_place() {
2081 use super::test_helpers::{make_label, test_db};
2082 let (db, _dir) = test_db();
2083
2084 let mut l = make_label("lbl_1", "Vanta", "default");
2085 db.upsert_label(&l).unwrap();
2086
2087 let listed = db.list_labels("default").unwrap();
2088 assert_eq!(listed.len(), 1);
2089 assert_eq!(listed[0].name, "Vanta");
2090
2091 l.name = "Compliance".to_string();
2093 db.upsert_label(&l).unwrap();
2094
2095 let listed = db.list_labels("default").unwrap();
2096 assert_eq!(listed.len(), 1);
2097 assert_eq!(listed[0].name, "Compliance");
2098 }
2099
2100 #[test]
2101 fn list_labels_is_workspace_scoped_and_sorted() {
2102 use super::test_helpers::{make_label, test_db};
2103 let (db, _dir) = test_db();
2104 db.upsert_workspace("work", None, None).unwrap();
2105
2106 db.upsert_label(&make_label("a", "Zebra", "default"))
2107 .unwrap();
2108 db.upsert_label(&make_label("b", "Apple", "default"))
2109 .unwrap();
2110 db.upsert_label(&make_label("c", "OnlyInWork", "work"))
2111 .unwrap();
2112
2113 let default_labels = db.list_labels("default").unwrap();
2114 assert_eq!(
2115 default_labels
2116 .iter()
2117 .map(|l| l.name.as_str())
2118 .collect::<Vec<_>>(),
2119 vec!["Apple", "Zebra"]
2120 );
2121 let work_labels = db.list_labels("work").unwrap();
2122 assert_eq!(work_labels.len(), 1);
2123 assert_eq!(work_labels[0].name, "OnlyInWork");
2124 }
2125
2126 #[test]
2127 fn delete_labels_for_workspace_not_in_removes_orphans() {
2128 use super::test_helpers::{make_label, test_db};
2129 let (db, _dir) = test_db();
2130
2131 db.upsert_label(&make_label("keep", "Keep", "default"))
2132 .unwrap();
2133 db.upsert_label(&make_label("drop", "Drop", "default"))
2134 .unwrap();
2135
2136 let kept = db
2137 .delete_labels_for_workspace_not_in("default", &["keep".to_string()])
2138 .unwrap();
2139 assert_eq!(kept, 1, "should report 1 deleted");
2140
2141 let listed = db.list_labels("default").unwrap();
2142 assert_eq!(listed.len(), 1);
2143 assert_eq!(listed[0].name, "Keep");
2144 }
2145
2146 #[test]
2147 fn replace_issue_labels_overwrites_existing() {
2148 use super::test_helpers::{make_issue, make_label, test_db};
2149 let (db, _dir) = test_db();
2150
2151 let issue = make_issue("ENG-1", "ENG");
2152 db.upsert_issue(&issue).unwrap();
2153 db.upsert_label(&make_label("l1", "Bug", "default"))
2154 .unwrap();
2155 db.upsert_label(&make_label("l2", "UI", "default")).unwrap();
2156 db.upsert_label(&make_label("l3", "Backend", "default"))
2157 .unwrap();
2158
2159 db.replace_issue_labels(&issue.id, &["l1".to_string(), "l2".to_string()])
2160 .unwrap();
2161 let labels = db.get_issue_label_ids(&issue.id).unwrap();
2162 assert_eq!(labels, vec!["l1".to_string(), "l2".to_string()]);
2163
2164 db.replace_issue_labels(&issue.id, &["l3".to_string()])
2166 .unwrap();
2167 let labels = db.get_issue_label_ids(&issue.id).unwrap();
2168 assert_eq!(labels, vec!["l3".to_string()]);
2169 }
2170
2171 #[test]
2172 fn deleting_issue_cascades_to_issue_labels() {
2173 use super::test_helpers::{make_issue, make_label, test_db};
2174 let (db, _dir) = test_db();
2175
2176 let issue = make_issue("ENG-2", "ENG");
2177 db.upsert_issue(&issue).unwrap();
2178 db.upsert_label(&make_label("l1", "Bug", "default"))
2179 .unwrap();
2180 db.replace_issue_labels(&issue.id, &["l1".to_string()])
2181 .unwrap();
2182
2183 db.with_conn(|conn| {
2184 conn.execute(
2185 "DELETE FROM issues WHERE id = ?1",
2186 rusqlite::params![&issue.id],
2187 )?;
2188 let n: i64 = conn.query_row(
2189 "SELECT COUNT(*) FROM issue_labels WHERE issue_id = ?1",
2190 rusqlite::params![&issue.id],
2191 |r| r.get(0),
2192 )?;
2193 assert_eq!(n, 0);
2194 Ok(())
2195 })
2196 .unwrap();
2197 }
2198
2199 #[test]
2200 fn deleting_label_cascades_to_issue_labels() {
2201 use super::test_helpers::{make_issue, make_label, test_db};
2202 let (db, _dir) = test_db();
2203
2204 let issue = make_issue("ENG-3", "ENG");
2205 db.upsert_issue(&issue).unwrap();
2206 db.upsert_label(&make_label("l1", "Bug", "default"))
2207 .unwrap();
2208 db.replace_issue_labels(&issue.id, &["l1".to_string()])
2209 .unwrap();
2210
2211 db.delete_labels_for_workspace_not_in("default", &[])
2212 .unwrap();
2213 let labels = db.get_issue_label_ids(&issue.id).unwrap();
2214 assert!(labels.is_empty());
2215 }
2216
2217 #[test]
2218 fn resolve_label_ids_local_matches_case_insensitive_and_returns_unknowns() {
2219 use super::test_helpers::{make_label, test_db};
2220 let (db, _dir) = test_db();
2221
2222 db.upsert_label(&make_label("l1", "Vanta", "default"))
2223 .unwrap();
2224 db.upsert_label(&make_label("l2", "Security", "default"))
2225 .unwrap();
2226
2227 let (resolved, unknown) = db
2228 .resolve_label_ids_local(
2229 "default",
2230 &[
2231 "vanta".to_string(),
2232 "secURity".to_string(),
2233 "missing".to_string(),
2234 ],
2235 )
2236 .unwrap();
2237 assert_eq!(resolved.len(), 2);
2238 assert!(resolved.contains(&"l1".to_string()));
2239 assert!(resolved.contains(&"l2".to_string()));
2240 assert_eq!(unknown, vec!["missing".to_string()]);
2241 }
2242
2243 #[test]
2244 fn resolve_label_ids_local_is_workspace_scoped() {
2245 use super::test_helpers::{make_label, test_db};
2246 let (db, _dir) = test_db();
2247 db.upsert_workspace("work", None, None).unwrap();
2248
2249 db.upsert_label(&make_label("l1", "Vanta", "default"))
2250 .unwrap();
2251 db.upsert_label(&make_label("l2", "Vanta", "work")).unwrap();
2252
2253 let (resolved, _) = db
2254 .resolve_label_ids_local("work", &["vanta".to_string()])
2255 .unwrap();
2256 assert_eq!(resolved, vec!["l2".to_string()]);
2257 }
2258
2259 #[test]
2260 fn get_unprioritized_issues_filters_by_labels_with_and_semantics() {
2261 use super::test_helpers::{make_issue, make_label, test_db};
2262 let (db, _dir) = test_db();
2263
2264 let mut a = make_issue("ENG-10", "ENG");
2265 a.priority = 0;
2266 let mut b = make_issue("ENG-11", "ENG");
2267 b.priority = 0;
2268 let mut c = make_issue("ENG-12", "ENG");
2269 c.priority = 0;
2270 db.upsert_issue(&a).unwrap();
2271 db.upsert_issue(&b).unwrap();
2272 db.upsert_issue(&c).unwrap();
2273
2274 db.upsert_label(&make_label("vanta", "Vanta", "default"))
2275 .unwrap();
2276 db.upsert_label(&make_label("sec", "Security", "default"))
2277 .unwrap();
2278
2279 db.replace_issue_labels(&a.id, &["vanta".to_string(), "sec".to_string()])
2280 .unwrap();
2281 db.replace_issue_labels(&b.id, &["vanta".to_string()])
2282 .unwrap();
2283 db.replace_issue_labels(&c.id, &["sec".to_string()])
2284 .unwrap();
2285
2286 let result = db
2288 .get_unprioritized_issues_filtered(
2289 Some("ENG"),
2290 false,
2291 "default",
2292 Some(&["vanta".to_string(), "sec".to_string()]),
2293 )
2294 .unwrap();
2295 let idents: Vec<_> = result.iter().map(|i| i.identifier.as_str()).collect();
2296 assert_eq!(idents, vec!["ENG-10"]);
2297
2298 let result = db
2300 .get_unprioritized_issues_filtered(
2301 Some("ENG"),
2302 false,
2303 "default",
2304 Some(&["vanta".to_string()]),
2305 )
2306 .unwrap();
2307 let idents: Vec<_> = result.iter().map(|i| i.identifier.as_str()).collect();
2308 assert!(idents.contains(&"ENG-10"));
2309 assert!(idents.contains(&"ENG-11"));
2310 assert!(!idents.contains(&"ENG-12"));
2311
2312 let result = db
2314 .get_unprioritized_issues_filtered(Some("ENG"), false, "default", None)
2315 .unwrap();
2316 assert_eq!(result.len(), 3);
2317 }
2318
2319 #[test]
2320 fn fts_search_with_label_filter_intersects() {
2321 use super::test_helpers::{make_issue, make_label, test_db};
2322 let (db, _dir) = test_db();
2323
2324 let mut a = make_issue("ENG-20", "ENG");
2325 a.title = "Audit logging gap".to_string();
2326 let mut b = make_issue("ENG-21", "ENG");
2327 b.title = "Audit something else".to_string();
2328 db.upsert_issue(&a).unwrap();
2329 db.upsert_issue(&b).unwrap();
2330
2331 db.upsert_label(&make_label("vanta", "Vanta", "default"))
2332 .unwrap();
2333 db.replace_issue_labels(&a.id, &["vanta".to_string()])
2334 .unwrap();
2335
2336 let r = db
2338 .fts_search_filtered("\"audit\"", 10, "default", None)
2339 .unwrap();
2340 assert_eq!(r.len(), 2);
2341
2342 let r = db
2344 .fts_search_filtered("\"audit\"", 10, "default", Some(&["vanta".to_string()]))
2345 .unwrap();
2346 assert_eq!(r.len(), 1);
2347 assert_eq!(r[0].identifier, "ENG-20");
2348 }
2349
2350 #[test]
2351 fn delete_workspace_cleans_up_labels() {
2352 use super::test_helpers::{make_label, test_db};
2353 let (db, _dir) = test_db();
2354 db.upsert_workspace("doomed", None, None).unwrap();
2355 db.upsert_label(&make_label("l1", "Lab", "doomed")).unwrap();
2356 db.delete_workspace("doomed").unwrap();
2357 let listed = db.list_labels("doomed").unwrap();
2358 assert!(listed.is_empty());
2359 }
2360
2361 #[test]
2362 fn embedding_selection_requires_hydrated_details_and_tracks_model_and_content() {
2363 let (db, _dir) = test_db();
2364 let mut issue = make_issue("CUT-1", "CUT");
2365 issue.id = "issue-1".into();
2366 db.upsert_issue_index(
2367 &super::IssueIndexEntry {
2368 id: issue.id.clone(),
2369 identifier: issue.identifier.clone(),
2370 team_key: issue.team_key.clone(),
2371 title: issue.title.clone(),
2372 state_name: issue.state_name.clone(),
2373 state_type: issue.state_type.clone(),
2374 created_at: issue.created_at.clone(),
2375 updated_at: issue.updated_at.clone(),
2376 archived_at: None,
2377 url: issue.url.clone(),
2378 },
2379 "default",
2380 "index-1",
2381 )
2382 .unwrap();
2383
2384 assert!(db
2385 .get_issues_needing_embedding_for_model(Some("CUT"), false, "default", Some("model-a"),)
2386 .unwrap()
2387 .is_empty());
2388
2389 db.mark_hydration_complete(
2390 "default",
2391 "issue-1",
2392 super::HydrationResource::Details,
2393 &issue.updated_at,
2394 "2026-01-03T00:00:00Z",
2395 )
2396 .unwrap();
2397 assert_eq!(
2398 db.get_issues_needing_embedding_for_model(
2399 Some("CUT"),
2400 false,
2401 "default",
2402 Some("model-a"),
2403 )
2404 .unwrap()
2405 .len(),
2406 1
2407 );
2408
2409 db.upsert_chunks_with_model(
2410 "issue-1",
2411 &[(0, "chunk".into(), fake_embedding(8))],
2412 "model-a",
2413 )
2414 .unwrap();
2415 assert!(db
2416 .get_issues_needing_embedding_for_model(Some("CUT"), false, "default", Some("model-a"),)
2417 .unwrap()
2418 .is_empty());
2419 assert_eq!(
2420 db.get_issues_needing_embedding_for_model(
2421 Some("CUT"),
2422 false,
2423 "default",
2424 Some("model-b"),
2425 )
2426 .unwrap()
2427 .len(),
2428 1
2429 );
2430
2431 let mut hydrated = db.get_issue("issue-1").unwrap().unwrap();
2432 hydrated.description = Some("new hydrated content".into());
2433 db.upsert_issue(&hydrated).unwrap();
2434 assert_eq!(db.count_embedded_issues(Some("CUT"), "default").unwrap(), 1);
2435 assert_eq!(
2436 db.get_issues_needing_embedding_for_model(
2437 Some("CUT"),
2438 false,
2439 "default",
2440 Some("model-a"),
2441 )
2442 .unwrap()
2443 .len(),
2444 1
2445 );
2446 }
2447
2448 #[test]
2449 fn index_change_waits_for_detail_rehydration_before_embedding() {
2450 let (db, _dir) = test_db();
2451 let mut issue = make_issue("CUT-1", "CUT");
2452 issue.id = "issue-1".into();
2453 db.upsert_issue(&issue).unwrap();
2454 db.ensure_hydration_state_for_issue("default", &issue, "migration")
2455 .unwrap();
2456 db.mark_hydration_complete(
2457 "default",
2458 &issue.id,
2459 super::HydrationResource::Details,
2460 &issue.updated_at,
2461 "2026-01-03T00:00:00Z",
2462 )
2463 .unwrap();
2464 db.upsert_chunks_with_model(
2465 &issue.id,
2466 &[(0, "chunk".into(), fake_embedding(8))],
2467 "model-a",
2468 )
2469 .unwrap();
2470
2471 let new_updated_at = "2026-01-04T00:00:00Z";
2472 db.upsert_issue_index(
2473 &super::IssueIndexEntry {
2474 id: issue.id.clone(),
2475 identifier: issue.identifier.clone(),
2476 team_key: issue.team_key.clone(),
2477 title: "new indexed title".into(),
2478 state_name: issue.state_name.clone(),
2479 state_type: issue.state_type.clone(),
2480 created_at: issue.created_at.clone(),
2481 updated_at: new_updated_at.into(),
2482 archived_at: None,
2483 url: issue.url.clone(),
2484 },
2485 "default",
2486 "index-2",
2487 )
2488 .unwrap();
2489 assert!(db
2490 .get_issues_needing_embedding_for_model(None, false, "default", Some("model-a"))
2491 .unwrap()
2492 .is_empty());
2493
2494 db.mark_hydration_complete(
2495 "default",
2496 &issue.id,
2497 super::HydrationResource::Details,
2498 new_updated_at,
2499 "2026-01-05T00:00:00Z",
2500 )
2501 .unwrap();
2502 assert_eq!(
2503 db.get_issues_needing_embedding_for_model(None, false, "default", Some("model-a"))
2504 .unwrap()
2505 .len(),
2506 1
2507 );
2508 }
2509
2510 #[test]
2511 fn migration_13_keeps_legacy_hydrated_issues_embedding_eligible() {
2512 let dir = tempfile::tempdir().unwrap();
2513 let path = dir.path().join("legacy-embedding.db");
2514 {
2515 let db = super::Database::open(&path).unwrap();
2516 let mut issue = make_issue("CUT-1", "CUT");
2517 issue.id = "issue-1".into();
2518 db.upsert_issue(&issue).unwrap();
2519 db.ensure_hydration_state_for_issue("default", &issue, "migration")
2520 .unwrap();
2521 db.mark_hydration_complete(
2522 "default",
2523 &issue.id,
2524 super::HydrationResource::Details,
2525 &issue.updated_at,
2526 "2026-01-03T00:00:00Z",
2527 )
2528 .unwrap();
2529 db.upsert_chunks_with_model(
2530 &issue.id,
2531 &[(0, "legacy chunk".into(), fake_embedding(8))],
2532 "model-a",
2533 )
2534 .unwrap();
2535 db.with_conn(|conn| {
2536 conn.execute("UPDATE chunks SET source_content_hash=''", [])?;
2537 conn.execute("DELETE FROM schema_version WHERE version=13", [])?;
2538 Ok(())
2539 })
2540 .unwrap();
2541 }
2542
2543 let migrated = super::Database::open(&path).unwrap();
2544 assert_eq!(
2545 migrated
2546 .get_issues_needing_embedding_for_model(
2547 Some("CUT"),
2548 false,
2549 "default",
2550 Some("model-a"),
2551 )
2552 .unwrap()
2553 .len(),
2554 1
2555 );
2556 }
2557}