smriti/db/
takeout_import_repo.rs1use std::collections::BTreeSet;
4
5use rusqlite::{params, Connection, OptionalExtension, Result as SqliteResult};
6
7pub struct TakeoutImportRepo<'a> {
8 conn: &'a Connection,
9}
10
11#[derive(Debug, Clone)]
12pub struct TakeoutLedgerItem {
13 pub content_hash: String,
14 pub file_path: String,
15 pub metadata_json: Option<String>,
16 pub albums: BTreeSet<String>,
17}
18
19impl<'a> TakeoutImportRepo<'a> {
20 pub fn new(conn: &'a Connection) -> Self {
21 Self { conn }
22 }
23
24 pub fn path_for_hash(&self, content_hash: &str) -> SqliteResult<Option<String>> {
25 self.conn
26 .query_row(
27 "SELECT file_path FROM google_takeout_items WHERE content_hash = ?1",
28 params![content_hash],
29 |row| row.get(0),
30 )
31 .optional()
32 }
33
34 pub fn upsert_item(
35 &self,
36 content_hash: &str,
37 file_path: &str,
38 metadata_json: Option<&str>,
39 albums: &BTreeSet<String>,
40 ) -> SqliteResult<()> {
41 let tx = self.conn.unchecked_transaction()?;
42 tx.execute(
43 r#"
44 INSERT INTO google_takeout_items (content_hash, file_path, metadata_json)
45 VALUES (?1, ?2, ?3)
46 ON CONFLICT(content_hash) DO UPDATE SET
47 file_path = excluded.file_path,
48 metadata_json = COALESCE(excluded.metadata_json, google_takeout_items.metadata_json),
49 updated_at = CURRENT_TIMESTAMP
50 "#,
51 params![content_hash, file_path, metadata_json],
52 )?;
53 for album in albums {
54 tx.execute(
55 "INSERT OR IGNORE INTO google_takeout_albums (content_hash, album_name) VALUES (?1, ?2)",
56 params![content_hash, album],
57 )?;
58 }
59 tx.commit()
60 }
61
62 pub fn upsert_items(&self, items: &[TakeoutLedgerItem]) -> SqliteResult<()> {
63 let tx = self.conn.unchecked_transaction()?;
64 for item in items {
65 tx.execute(
66 r#"
67 INSERT INTO google_takeout_items (content_hash, file_path, metadata_json)
68 VALUES (?1, ?2, ?3)
69 ON CONFLICT(content_hash) DO UPDATE SET
70 file_path = excluded.file_path,
71 metadata_json = COALESCE(excluded.metadata_json, google_takeout_items.metadata_json),
72 updated_at = CURRENT_TIMESTAMP
73 "#,
74 params![item.content_hash, item.file_path, item.metadata_json],
75 )?;
76 for album in &item.albums {
77 tx.execute(
78 "INSERT OR IGNORE INTO google_takeout_albums (content_hash, album_name) VALUES (?1, ?2)",
79 params![item.content_hash, album],
80 )?;
81 }
82 }
83 tx.commit()
84 }
85
86 pub fn candidate_existing_files(
87 &self,
88 sizes: &std::collections::HashSet<i64>,
89 ) -> SqliteResult<Vec<(String, i64)>> {
90 let mut stmt = self
91 .conn
92 .prepare("SELECT file_path, file_size FROM photos WHERE is_trashed = FALSE")?;
93 let rows = stmt.query_map([], |row| Ok((row.get(0)?, row.get(1)?)))?;
94 let mut result = Vec::new();
95 for row in rows {
96 let row = row?;
97 if sizes.contains(&row.1) {
98 result.push(row);
99 }
100 }
101 Ok(result)
102 }
103}
104
105#[cfg(test)]
106mod tests {
107 use super::*;
108 use crate::db::Database;
109 use tempfile::tempdir;
110
111 #[test]
112 fn upsert_is_idempotent_and_unions_albums() {
113 let dir = tempdir().unwrap();
114 let db = Database::open_for_drive(dir.path()).unwrap();
115 crate::db::create_schema(&db.conn).unwrap();
116 let repo = TakeoutImportRepo::new(&db.conn);
117 repo.upsert_item(
118 "abc",
119 "Imported from Google Photos/2020/a.jpg",
120 Some(r#"{"favorited":true}"#),
121 &BTreeSet::from(["Trip".to_string()]),
122 )
123 .unwrap();
124 repo.upsert_item(
125 "abc",
126 "Imported from Google Photos/2020/a.jpg",
127 None,
128 &BTreeSet::from(["Family".to_string()]),
129 )
130 .unwrap();
131
132 assert_eq!(
133 repo.path_for_hash("abc").unwrap().as_deref(),
134 Some("Imported from Google Photos/2020/a.jpg")
135 );
136 let count: i64 = db
137 .conn
138 .query_row(
139 "SELECT COUNT(*) FROM google_takeout_albums WHERE content_hash = 'abc'",
140 [],
141 |r| r.get(0),
142 )
143 .unwrap();
144 assert_eq!(count, 2);
145 }
146}