Skip to main content

max / audiofiles

47.0 KB · 1354 lines History Blame Raw
1 //! State tracking: snapshots, changelog maintenance, cloud-only marking.
2
3 use std::path::Path;
4
5 use rusqlite::Connection;
6
7 use tracing::instrument;
8
9 use crate::error::Result;
10
11 use super::{get_sync_state, set_sync_state, MAX_CHANGELOG_ENTRIES};
12
13 /// Create initial snapshot: insert all existing rows into sync_changelog.
14 #[instrument(skip_all)]
15 pub fn create_initial_snapshot(conn: &Connection) -> Result<i64> {
16 let done = get_sync_state(conn, "initial_snapshot_done")?;
17 if done == "1" {
18 return Ok(0);
19 }
20
21 let tx = conn.unchecked_transaction()?;
22
23 let mut total: i64 = 0;
24
25 let table_queries: &[(&str, &str)] = &[
26 ("samples", "SELECT hash, json_object('hash', hash, 'original_name', original_name, 'file_extension', file_extension, 'file_size', file_size, 'import_date', import_date, 'last_modified', last_modified, 'cloud_only', cloud_only, 'duration', duration) FROM samples"),
27 ("audio_analysis", "SELECT hash, json_object('hash', hash, 'bpm', bpm, 'musical_key', musical_key, 'duration', duration, 'sample_rate', sample_rate, 'channels', channels, 'peak_db', peak_db, 'rms_db', rms_db, 'is_loop', is_loop, 'spectral_centroid', spectral_centroid, 'onset_strength', onset_strength, 'analyzed_at', analyzed_at, 'lufs', lufs, 'spectral_flatness', spectral_flatness, 'spectral_rolloff', spectral_rolloff, 'zero_crossing_rate', zero_crossing_rate, 'classification', classification, 'spectral_bandwidth', spectral_bandwidth, 'centroid_variance', centroid_variance, 'crest_factor', crest_factor, 'attack_time', attack_time, 'classification_confidence', classification_confidence) FROM audio_analysis"),
28 ("vfs", "SELECT CAST(id AS TEXT), json_object('id', id, 'name', name, 'created_at', created_at, 'modified_at', modified_at, 'sync_files', sync_files) FROM vfs"),
29 ("vfs_nodes", "SELECT CAST(id AS TEXT), json_object('id', id, 'vfs_id', vfs_id, 'parent_id', parent_id, 'name', name, 'node_type', node_type, 'sample_hash', sample_hash, 'created_at', created_at) FROM vfs_nodes"),
30 ("tags", "SELECT sample_hash || ':' || tag, json_object('sample_hash', sample_hash, 'tag', tag) FROM tags"),
31 ("collections", "SELECT CAST(id AS TEXT), json_object('id', id, 'name', name, 'description', description, 'created_at', created_at, 'filter_json', filter_json) FROM collections"),
32 ("collection_members", "SELECT CAST(collection_id AS TEXT) || ':' || sample_hash, json_object('collection_id', collection_id, 'sample_hash', sample_hash, 'added_at', added_at) FROM collection_members"),
33 ("user_config", "SELECT key, json_object('key', key, 'value', value) FROM user_config WHERE key NOT LIKE 'sync_%' AND key != 'loose_files'"),
34 ("edit_history", "SELECT CAST(id AS TEXT), json_object('id', id, 'source_hash', source_hash, 'result_hash', result_hash, 'operation', operation, 'params_json', params_json, 'created_at', created_at) FROM edit_history"),
35 ];
36
37 for (table, query) in table_queries {
38 let mut stmt = tx.prepare(query)?;
39 let rows: Vec<(String, String)> = stmt
40 .query_map([], |row| Ok((row.get(0)?, row.get(1)?)))?
41 .collect::<std::result::Result<Vec<_>, _>>()?;
42
43 for (row_id, data) in &rows {
44 tx.execute(
45 "INSERT INTO sync_changelog (table_name, op, row_id, data) VALUES (?1, 'INSERT', ?2, ?3)",
46 rusqlite::params![table, row_id, data],
47 )?;
48 }
49 total += rows.len() as i64;
50 }
51
52 set_sync_state(&tx, "initial_snapshot_done", "1")?;
53 tx.commit()?;
54 Ok(total)
55 }
56
57 /// Delete pushed changelog entries older than 7 days.
58 #[instrument(skip_all)]
59 pub fn cleanup_changelog(conn: &Connection) -> Result<i64> {
60 let deleted = conn.execute(
61 "DELETE FROM sync_changelog WHERE pushed = 1 AND timestamp < datetime('now', '-7 days')",
62 [],
63 )?;
64 Ok(deleted as i64)
65 }
66
67 /// Enforce a hard cap on total changelog entries to prevent unbounded growth
68 /// when sync is disconnected or failing. Deletes the oldest entries (by rowid)
69 /// to bring the count back under `MAX_CHANGELOG_ENTRIES`.
70 #[instrument(skip_all)]
71 pub fn enforce_changelog_retention(conn: &Connection) -> Result<i64> {
72 let count: i64 = conn.query_row("SELECT COUNT(*) FROM sync_changelog", [], |r| r.get(0))?;
73
74 if count <= MAX_CHANGELOG_ENTRIES {
75 return Ok(0);
76 }
77
78 let excess = count - MAX_CHANGELOG_ENTRIES;
79
80 // Prefer deleting already-pushed entries first to avoid losing unsynced changes.
81 let deleted_pushed = conn.execute(
82 "DELETE FROM sync_changelog WHERE id IN \
83 (SELECT id FROM sync_changelog WHERE pushed = 1 ORDER BY id ASC LIMIT ?1)",
84 [excess],
85 )?;
86
87 let mut total_deleted = deleted_pushed as i64;
88
89 // If still over cap, reluctantly delete unpushed entries (oldest first).
90 if total_deleted < excess {
91 let remaining = excess - total_deleted;
92 let deleted_unpushed = conn.execute(
93 "DELETE FROM sync_changelog WHERE id IN \
94 (SELECT id FROM sync_changelog ORDER BY id ASC LIMIT ?1)",
95 [remaining],
96 )?;
97 if deleted_unpushed > 0 {
98 tracing::warn!(
99 deleted_unpushed,
100 "Changelog retention: dropped unpushed entries — sync was offline too long"
101 );
102 }
103 total_deleted += deleted_unpushed as i64;
104 }
105
106 tracing::warn!(
107 deleted = total_deleted,
108 total = count,
109 limit = MAX_CHANGELOG_ENTRIES,
110 "Changelog retention cap enforced"
111 );
112 Ok(total_deleted)
113 }
114
115 /// Mark samples as cloud_only when their blob doesn't exist on disk.
116 ///
117 /// After pulling remote changes, a sample row may be created with cloud_only=0
118 /// (as it was on the origin device). If the local content directory doesn't have
119 /// the blob file, we correct the flag to cloud_only=1. This runs with
120 /// applying_remote=1 to suppress changelog triggers.
121 #[instrument(skip_all)]
122 pub fn mark_cloud_only_samples(conn: &Connection, content_dir: &Path) -> Result<i64> {
123 let mut stmt = conn.prepare(
124 "SELECT hash, file_extension FROM samples WHERE cloud_only = 0",
125 )?;
126 let rows: Vec<(String, String)> = stmt
127 .query_map([], |row| Ok((row.get(0)?, row.get(1)?)))?
128 .collect::<std::result::Result<Vec<_>, _>>()?;
129
130 // Collect hashes missing on disk before taking a transaction.
131 let missing: Vec<&str> = rows
132 .iter()
133 .filter(|(hash, ext)| !content_dir.join(format!("{}.{}", hash, ext)).exists())
134 .map(|(hash, _)| hash.as_str())
135 .collect();
136
137 if missing.is_empty() {
138 return Ok(0);
139 }
140
141 conn.execute_batch("BEGIN IMMEDIATE")?;
142 conn.execute(
143 "UPDATE sync_state SET value = '1' WHERE key = 'applying_remote'",
144 [],
145 )?;
146 let result = (|| -> Result<i64> {
147 let mut marked = 0i64;
148 for hash in &missing {
149 conn.execute(
150 "UPDATE samples SET cloud_only = 1 WHERE hash = ?1",
151 [hash],
152 )?;
153 marked += 1;
154 }
155 Ok(marked)
156 })();
157 conn.execute(
158 "UPDATE sync_state SET value = '0' WHERE key = 'applying_remote'",
159 [],
160 )?;
161 match result {
162 Ok(marked) => {
163 conn.execute_batch("COMMIT")?;
164 Ok(marked)
165 }
166 Err(e) => {
167 let _ = conn.execute_batch("ROLLBACK");
168 Err(e)
169 }
170 }
171 }
172
173 #[cfg(test)]
174 mod tests {
175 use super::*;
176 use super::super::{
177 UPSERT_ORDER, DELETE_ORDER, table_columns, pk_columns,
178 resolve::{apply_upsert, apply_delete, apply_remote_changes},
179 };
180 use audiofiles_core::db::Database;
181 use serde_json::json;
182 use synckit_client::{ChangeEntry, ChangeOp};
183
184 fn setup_test_db() -> Database {
185 Database::open_in_memory().expect("Failed to create test DB")
186 }
187
188 fn insert_sample(conn: &Connection, hash: &str, name: &str, ext: &str) {
189 let now = chrono::Utc::now().timestamp();
190 conn.execute(
191 "INSERT INTO samples (hash, original_name, file_extension, file_size, import_date, last_modified) VALUES (?1, ?2, ?3, 1024, ?4, ?4)",
192 rusqlite::params![hash, name, ext, now],
193 ).unwrap();
194 }
195
196 fn insert_vfs(conn: &Connection, name: &str, sync_files: bool) -> i64 {
197 let now = chrono::Utc::now().timestamp();
198 conn.execute(
199 "INSERT INTO vfs (name, created_at, modified_at, sync_files) VALUES (?1, ?2, ?2, ?3)",
200 rusqlite::params![name, now, sync_files as i64],
201 ).unwrap();
202 conn.last_insert_rowid()
203 }
204
205 fn clear_changelog(conn: &Connection) {
206 conn.execute("DELETE FROM sync_changelog", []).unwrap();
207 }
208
209 fn changelog_count(conn: &Connection, table: Option<&str>, op: Option<&str>) -> i64 {
210 match (table, op) {
211 (Some(t), Some(o)) => conn.query_row(
212 "SELECT COUNT(*) FROM sync_changelog WHERE table_name = ?1 AND op = ?2",
213 rusqlite::params![t, o],
214 |row| row.get(0),
215 ).unwrap(),
216 (Some(t), None) => conn.query_row(
217 "SELECT COUNT(*) FROM sync_changelog WHERE table_name = ?1",
218 [t],
219 |row| row.get(0),
220 ).unwrap(),
221 (None, Some(o)) => conn.query_row(
222 "SELECT COUNT(*) FROM sync_changelog WHERE op = ?1",
223 [o],
224 |row| row.get(0),
225 ).unwrap(),
226 (None, None) => conn.query_row(
227 "SELECT COUNT(*) FROM sync_changelog",
228 [],
229 |row| row.get(0),
230 ).unwrap(),
231 }
232 }
233
234 fn change(table: &str, op: ChangeOp, row_id: &str, data: Option<serde_json::Value>) -> ChangeEntry {
235 ChangeEntry {
236 table: table.to_string(),
237 op,
238 row_id: row_id.to_string(),
239 timestamp: chrono::Utc::now(),
240 data,
241 }
242 }
243
244 // ── FK ordering ──
245
246 #[test]
247 fn upsert_order_parents_before_children() {
248 let pos = |t: &str| UPSERT_ORDER.iter().position(|x| *x == t).unwrap();
249 assert!(pos("vfs") < pos("vfs_nodes"));
250 assert!(pos("samples") < pos("audio_analysis"));
251 assert!(pos("samples") < pos("tags"));
252 assert!(pos("samples") < pos("collection_members"));
253 assert!(pos("collections") < pos("collection_members"));
254 }
255
256 #[test]
257 fn delete_order_children_before_parents() {
258 let pos = |t: &str| DELETE_ORDER.iter().position(|x| *x == t).unwrap();
259 assert!(pos("vfs_nodes") < pos("vfs"));
260 assert!(pos("audio_analysis") < pos("samples"));
261 assert!(pos("tags") < pos("samples"));
262 assert!(pos("collection_members") < pos("collections"));
263 assert!(pos("collection_members") < pos("samples"));
264 }
265
266 #[test]
267 fn orders_are_exact_reverses() {
268 let reversed: Vec<&str> = UPSERT_ORDER.iter().rev().copied().collect();
269 assert_eq!(reversed, DELETE_ORDER);
270 }
271
272 // ── Column whitelists ──
273
274 #[test]
275 fn all_tables_have_column_whitelists() {
276 for table in UPSERT_ORDER {
277 assert!(
278 table_columns(table).is_some(),
279 "missing column whitelist for: {}", table
280 );
281 }
282 }
283
284 #[test]
285 fn unknown_table_returns_none() {
286 assert!(table_columns("nonexistent").is_none());
287 assert!(table_columns("fingerprints").is_none());
288 }
289
290 #[test]
291 fn pk_columns_covers_all_tables() {
292 for table in UPSERT_ORDER {
293 let pks = pk_columns(table);
294 assert!(!pks.is_empty(), "missing pk_columns for: {}", table);
295 }
296 assert_eq!(pk_columns("tags"), &["sample_hash", "tag"]);
297 assert_eq!(pk_columns("collection_members"), &["collection_id", "sample_hash"]);
298 }
299
300 // ── Triggers ──
301
302 #[test]
303 fn sample_insert_fires_trigger() {
304 let db = setup_test_db();
305 let conn = db.conn();
306 clear_changelog(conn);
307
308 insert_sample(conn, "abc123", "kick.wav", "wav");
309
310 assert_eq!(changelog_count(conn, Some("samples"), Some("INSERT")), 1);
311
312 let data: String = conn.query_row(
313 "SELECT data FROM sync_changelog WHERE table_name = 'samples' AND row_id = 'abc123'",
314 [],
315 |row| row.get(0),
316 ).unwrap();
317 let parsed: serde_json::Value = serde_json::from_str(&data).unwrap();
318 assert!(parsed.get("cloud_only").is_some());
319 assert_eq!(parsed["hash"], "abc123");
320 }
321
322 #[test]
323 fn vfs_insert_fires_trigger() {
324 let db = setup_test_db();
325 let conn = db.conn();
326 clear_changelog(conn);
327
328 let vfs_id = insert_vfs(conn, "Library", true);
329
330 assert_eq!(changelog_count(conn, Some("vfs"), Some("INSERT")), 1);
331
332 let row_id: String = conn.query_row(
333 "SELECT row_id FROM sync_changelog WHERE table_name = 'vfs'",
334 [],
335 |row| row.get(0),
336 ).unwrap();
337 assert_eq!(row_id, vfs_id.to_string());
338 }
339
340 #[test]
341 fn tag_insert_fires_trigger() {
342 let db = setup_test_db();
343 let conn = db.conn();
344 insert_sample(conn, "hash1", "snare.wav", "wav");
345 clear_changelog(conn);
346
347 conn.execute(
348 "INSERT INTO tags (sample_hash, tag) VALUES ('hash1', 'drums')",
349 [],
350 ).unwrap();
351
352 assert_eq!(changelog_count(conn, Some("tags"), Some("INSERT")), 1);
353
354 let row_id: String = conn.query_row(
355 "SELECT row_id FROM sync_changelog WHERE table_name = 'tags'",
356 [],
357 |row| row.get(0),
358 ).unwrap();
359 assert_eq!(row_id, "hash1:drums");
360 }
361
362 #[test]
363 fn collection_member_insert_fires_trigger() {
364 let db = setup_test_db();
365 let conn = db.conn();
366 insert_sample(conn, "hash2", "hat.wav", "wav");
367 let now = chrono::Utc::now().timestamp();
368 conn.execute(
369 "INSERT INTO collections (name, description, created_at) VALUES ('Faves', NULL, ?1)",
370 [now],
371 ).unwrap();
372 let collection_id = conn.last_insert_rowid();
373 clear_changelog(conn);
374
375 conn.execute(
376 "INSERT INTO collection_members (collection_id, sample_hash, added_at) VALUES (?1, 'hash2', ?2)",
377 rusqlite::params![collection_id, now],
378 ).unwrap();
379
380 assert_eq!(changelog_count(conn, Some("collection_members"), Some("INSERT")), 1);
381
382 let row_id: String = conn.query_row(
383 "SELECT row_id FROM sync_changelog WHERE table_name = 'collection_members'",
384 [],
385 |row| row.get(0),
386 ).unwrap();
387 assert_eq!(row_id, format!("{}:hash2", collection_id));
388 }
389
390 // ── Trigger suppression ──
391
392 #[test]
393 fn trigger_suppression_during_remote_apply() {
394 let db = setup_test_db();
395 let conn = db.conn();
396 clear_changelog(conn);
397
398 set_sync_state(conn, "applying_remote", "1").unwrap();
399 insert_sample(conn, "suppressed", "test.wav", "wav");
400 assert_eq!(changelog_count(conn, None, None), 0);
401
402 set_sync_state(conn, "applying_remote", "0").unwrap();
403 insert_sample(conn, "unsuppressed", "test2.wav", "wav");
404 assert_eq!(changelog_count(conn, None, None), 1);
405 }
406
407 // ── apply_upsert ──
408
409 #[test]
410 fn apply_upsert_inserts_sample() {
411 let db = setup_test_db();
412 let conn = db.conn();
413
414 let data = json!({
415 "hash": "upsert_hash",
416 "original_name": "synced.wav",
417 "file_extension": "wav",
418 "file_size": 2048,
419 "import_date": 1000000,
420 "last_modified": 1000000,
421 "cloud_only": 0
422 });
423
424 apply_upsert(conn, "samples", &data).unwrap();
425
426 let name: String = conn.query_row(
427 "SELECT original_name FROM samples WHERE hash = 'upsert_hash'",
428 [],
429 |row| row.get(0),
430 ).unwrap();
431 assert_eq!(name, "synced.wav");
432 }
433
434 #[test]
435 fn apply_upsert_inserts_vfs_node() {
436 let db = setup_test_db();
437 let conn = db.conn();
438
439 let vfs_id = insert_vfs(conn, "TestVFS", false);
440 insert_sample(conn, "node_hash", "pad.wav", "wav");
441
442 let data = json!({
443 "id": 999,
444 "vfs_id": vfs_id,
445 "parent_id": null,
446 "name": "pad.wav",
447 "node_type": "sample",
448 "sample_hash": "node_hash",
449 "created_at": 1000000
450 });
451
452 apply_upsert(conn, "vfs_nodes", &data).unwrap();
453
454 let name: String = conn.query_row(
455 "SELECT name FROM vfs_nodes WHERE id = 999",
456 [],
457 |row| row.get(0),
458 ).unwrap();
459 assert_eq!(name, "pad.wav");
460 }
461
462 #[test]
463 fn apply_upsert_unknown_table_is_no_op() {
464 let db = setup_test_db();
465 let conn = db.conn();
466
467 let data = json!({"id": "abc"});
468 let result = apply_upsert(conn, "nonexistent_table", &data);
469 assert!(result.is_ok());
470 }
471
472 // ── apply_delete ──
473
474 #[test]
475 fn apply_delete_removes_sample() {
476 let db = setup_test_db();
477 let conn = db.conn();
478 insert_sample(conn, "del_hash", "delete_me.wav", "wav");
479
480 apply_delete(conn, "samples", "del_hash").unwrap();
481
482 let count: i64 = conn.query_row(
483 "SELECT COUNT(*) FROM samples WHERE hash = 'del_hash'",
484 [],
485 |row| row.get(0),
486 ).unwrap();
487 assert_eq!(count, 0);
488 }
489
490 #[test]
491 fn apply_delete_composite_pk_tags() {
492 let db = setup_test_db();
493 let conn = db.conn();
494 insert_sample(conn, "tag_hash", "tagged.wav", "wav");
495 conn.execute(
496 "INSERT INTO tags (sample_hash, tag) VALUES ('tag_hash', 'bass')",
497 [],
498 ).unwrap();
499
500 apply_delete(conn, "tags", "tag_hash:bass").unwrap();
501
502 let count: i64 = conn.query_row(
503 "SELECT COUNT(*) FROM tags WHERE sample_hash = 'tag_hash' AND tag = 'bass'",
504 [],
505 |row| row.get(0),
506 ).unwrap();
507 assert_eq!(count, 0);
508 }
509
510 // ── Full pipeline ──
511
512 #[test]
513 fn apply_remote_changes_full_pipeline() {
514 let db = setup_test_db();
515 let conn = db.conn();
516 clear_changelog(conn);
517
518 // Changes in wrong FK order — apply_remote_changes reorders via UPSERT_ORDER
519 let changes = vec![
520 change("vfs_nodes", ChangeOp::Insert, "1", Some(json!({
521 "id": 1, "vfs_id": 1, "parent_id": null,
522 "name": "kick.wav", "node_type": "sample",
523 "sample_hash": "pipe_hash", "created_at": 1000000
524 }))),
525 change("vfs", ChangeOp::Insert, "1", Some(json!({
526 "id": 1, "name": "Library", "created_at": 1000000,
527 "modified_at": 1000000, "sync_files": 1
528 }))),
529 change("samples", ChangeOp::Insert, "pipe_hash", Some(json!({
530 "hash": "pipe_hash", "original_name": "kick.wav",
531 "file_extension": "wav", "file_size": 4096,
532 "import_date": 1000000, "last_modified": 1000000,
533 "cloud_only": 0
534 }))),
535 ];
536
537 let applied = apply_remote_changes(conn, &changes).unwrap();
538 assert_eq!(applied, 3);
539
540 let sample_count: i64 = conn.query_row(
541 "SELECT COUNT(*) FROM samples WHERE hash = 'pipe_hash'",
542 [], |row| row.get(0),
543 ).unwrap();
544 assert_eq!(sample_count, 1);
545
546 let vfs_count: i64 = conn.query_row(
547 "SELECT COUNT(*) FROM vfs WHERE name = 'Library'",
548 [], |row| row.get(0),
549 ).unwrap();
550 assert_eq!(vfs_count, 1);
551
552 let node_count: i64 = conn.query_row(
553 "SELECT COUNT(*) FROM vfs_nodes WHERE sample_hash = 'pipe_hash'",
554 [], |row| row.get(0),
555 ).unwrap();
556 assert_eq!(node_count, 1);
557
558 // Changelog should be empty (triggers were suppressed)
559 assert_eq!(changelog_count(conn, None, None), 0);
560 }
561
562 // ── Initial snapshot ──
563
564 #[test]
565 fn create_initial_snapshot_captures_all_rows() {
566 let db = setup_test_db();
567 let conn = db.conn();
568
569 // Suppress triggers during data setup
570 set_sync_state(conn, "applying_remote", "1").unwrap();
571 insert_sample(conn, "snap1", "one.wav", "wav");
572 insert_sample(conn, "snap2", "two.wav", "wav");
573 insert_vfs(conn, "TestLib", false);
574 set_sync_state(conn, "applying_remote", "0").unwrap();
575
576 clear_changelog(conn);
577
578 let total = create_initial_snapshot(conn).unwrap();
579 assert_eq!(total, 3); // 2 samples + 1 vfs
580 }
581
582 #[test]
583 fn create_initial_snapshot_idempotent() {
584 let db = setup_test_db();
585 let conn = db.conn();
586
587 set_sync_state(conn, "applying_remote", "1").unwrap();
588 insert_sample(conn, "idem", "test.wav", "wav");
589 set_sync_state(conn, "applying_remote", "0").unwrap();
590 clear_changelog(conn);
591
592 let first = create_initial_snapshot(conn).unwrap();
593 assert!(first > 0);
594
595 let second = create_initial_snapshot(conn).unwrap();
596 assert_eq!(second, 0);
597 }
598
599 // ── Changelog helpers ──
600
601 #[test]
602 fn cleanup_removes_old_pushed() {
603 let db = setup_test_db();
604 let conn = db.conn();
605 clear_changelog(conn);
606
607 // Old pushed — should be deleted
608 conn.execute(
609 "INSERT INTO sync_changelog (table_name, op, row_id, timestamp, pushed) \
610 VALUES ('samples', 'INSERT', 'old', datetime('now', '-10 days'), 1)",
611 [],
612 ).unwrap();
613
614 // Recent pushed — should remain
615 conn.execute(
616 "INSERT INTO sync_changelog (table_name, op, row_id, timestamp, pushed) \
617 VALUES ('samples', 'INSERT', 'recent', datetime('now'), 1)",
618 [],
619 ).unwrap();
620
621 // Old unpushed — should remain (never delete unpushed)
622 conn.execute(
623 "INSERT INTO sync_changelog (table_name, op, row_id, timestamp, pushed) \
624 VALUES ('samples', 'INSERT', 'unpushed', datetime('now', '-10 days'), 0)",
625 [],
626 ).unwrap();
627
628 let deleted = cleanup_changelog(conn).unwrap();
629 assert_eq!(deleted, 1);
630
631 let remaining: i64 = conn.query_row(
632 "SELECT COUNT(*) FROM sync_changelog",
633 [], |row| row.get(0),
634 ).unwrap();
635 assert_eq!(remaining, 2);
636 }
637
638 #[test]
639 fn enforce_retention_caps_entries() {
640 let db = setup_test_db();
641 let conn = db.conn();
642 clear_changelog(conn);
643
644 let total = MAX_CHANGELOG_ENTRIES + 500;
645 for i in 0..total {
646 conn.execute(
647 "INSERT INTO sync_changelog (table_name, op, row_id, data) \
648 VALUES ('samples', 'UPDATE', ?1, '{}')",
649 [format!("row-{}", i)],
650 )
651 .unwrap();
652 }
653
654 let before: i64 =
655 conn.query_row("SELECT COUNT(*) FROM sync_changelog", [], |r| r.get(0)).unwrap();
656 assert_eq!(before, total);
657
658 let deleted = enforce_changelog_retention(conn).unwrap();
659 assert_eq!(deleted, 500);
660
661 let after: i64 =
662 conn.query_row("SELECT COUNT(*) FROM sync_changelog", [], |r| r.get(0)).unwrap();
663 assert_eq!(after, MAX_CHANGELOG_ENTRIES);
664
665 // Oldest entries removed — the remaining start at row-500
666 let min_row_id: String = conn
667 .query_row(
668 "SELECT row_id FROM sync_changelog ORDER BY id ASC LIMIT 1",
669 [],
670 |r| r.get(0),
671 )
672 .unwrap();
673 assert_eq!(min_row_id, "row-500");
674 }
675
676 #[test]
677 fn enforce_retention_noop_under_cap() {
678 let db = setup_test_db();
679 let conn = db.conn();
680 clear_changelog(conn);
681
682 for i in 0..5 {
683 conn.execute(
684 "INSERT INTO sync_changelog (table_name, op, row_id, data) \
685 VALUES ('samples', 'INSERT', ?1, '{}')",
686 [format!("row-{}", i)],
687 )
688 .unwrap();
689 }
690
691 let deleted = enforce_changelog_retention(conn).unwrap();
692 assert_eq!(deleted, 0);
693 }
694
695 #[test]
696 fn count_pending_counts_unpushed() {
697 let db = setup_test_db();
698 let conn = db.conn();
699 clear_changelog(conn);
700
701 // 3 unpushed
702 for i in 0..3 {
703 conn.execute(
704 "INSERT INTO sync_changelog (table_name, op, row_id, pushed) VALUES ('samples', 'INSERT', ?1, 0)",
705 [format!("un-{}", i)],
706 ).unwrap();
707 }
708
709 // 2 pushed
710 for i in 0..2 {
711 conn.execute(
712 "INSERT INTO sync_changelog (table_name, op, row_id, pushed) VALUES ('samples', 'INSERT', ?1, 1)",
713 [format!("push-{}", i)],
714 ).unwrap();
715 }
716
717 let count = super::super::count_pending_changes(conn).unwrap();
718 assert_eq!(count, 3);
719 }
720
721 // --- mark_cloud_only_samples tests ---
722
723 #[test]
724 fn mark_cloud_only_marks_missing_blobs() {
725 let db = setup_test_db();
726 let conn = db.conn();
727
728 // Insert two samples — neither has a file on disk
729 insert_sample(conn, "aaa", "kick.wav", "wav");
730 insert_sample(conn, "bbb", "snare.wav", "wav");
731
732 // Use a temp dir as content_dir (empty — no blobs)
733 let tmp = tempfile::tempdir().unwrap();
734
735 let marked = mark_cloud_only_samples(conn, tmp.path()).unwrap();
736 assert_eq!(marked, 2);
737
738 // Verify both are now cloud_only=1
739 let count: i64 = conn.query_row(
740 "SELECT COUNT(*) FROM samples WHERE cloud_only = 1",
741 [],
742 |r| r.get(0),
743 ).unwrap();
744 assert_eq!(count, 2);
745 }
746
747 #[test]
748 fn mark_cloud_only_skips_existing_blobs() {
749 let db = setup_test_db();
750 let conn = db.conn();
751
752 insert_sample(conn, "aaa", "kick.wav", "wav");
753 insert_sample(conn, "bbb", "snare.wav", "wav");
754
755 // Create one blob file, leave other missing
756 let tmp = tempfile::tempdir().unwrap();
757 std::fs::write(tmp.path().join("aaa.wav"), b"fake audio data").unwrap();
758
759 let marked = mark_cloud_only_samples(conn, tmp.path()).unwrap();
760 assert_eq!(marked, 1); // only bbb
761
762 // aaa still cloud_only=0, bbb is cloud_only=1
763 let aaa_co: i32 = conn.query_row(
764 "SELECT cloud_only FROM samples WHERE hash = 'aaa'",
765 [],
766 |r| r.get(0),
767 ).unwrap();
768 assert_eq!(aaa_co, 0);
769
770 let bbb_co: i32 = conn.query_row(
771 "SELECT cloud_only FROM samples WHERE hash = 'bbb'",
772 [],
773 |r| r.get(0),
774 ).unwrap();
775 assert_eq!(bbb_co, 1);
776 }
777
778 #[test]
779 fn mark_cloud_only_suppresses_changelog() {
780 let db = setup_test_db();
781 let conn = db.conn();
782
783 insert_sample(conn, "aaa", "kick.wav", "wav");
784 clear_changelog(conn);
785
786 let tmp = tempfile::tempdir().unwrap();
787 mark_cloud_only_samples(conn, tmp.path()).unwrap();
788
789 // The UPDATE should not appear in changelog (applying_remote suppression)
790 let count = changelog_count(conn, Some("samples"), Some("UPDATE"));
791 assert_eq!(count, 0);
792 }
793
794 #[test]
795 fn mark_cloud_only_idempotent() {
796 let db = setup_test_db();
797 let conn = db.conn();
798
799 insert_sample(conn, "aaa", "kick.wav", "wav");
800 let tmp = tempfile::tempdir().unwrap();
801
802 // First call marks it
803 assert_eq!(mark_cloud_only_samples(conn, tmp.path()).unwrap(), 1);
804 // Second call: already cloud_only=1, query only selects cloud_only=0
805 assert_eq!(mark_cloud_only_samples(conn, tmp.path()).unwrap(), 0);
806 }
807
808 // ── Integration: multi-table snapshot ──
809
810 #[test]
811 fn snapshot_captures_tags_collections_and_members() {
812 let db = setup_test_db();
813 let conn = db.conn();
814
815 // Suppress triggers during data setup
816 set_sync_state(conn, "applying_remote", "1").unwrap();
817
818 insert_sample(conn, "s1", "kick.wav", "wav");
819 insert_sample(conn, "s2", "snare.wav", "wav");
820
821 // Tags
822 conn.execute(
823 "INSERT INTO tags (sample_hash, tag) VALUES ('s1', 'drums')",
824 [],
825 ).unwrap();
826 conn.execute(
827 "INSERT INTO tags (sample_hash, tag) VALUES ('s2', 'perc')",
828 [],
829 ).unwrap();
830
831 // Collection
832 let now = chrono::Utc::now().timestamp();
833 conn.execute(
834 "INSERT INTO collections (name, description, created_at) VALUES ('Kit', NULL, ?1)",
835 [now],
836 ).unwrap();
837 let coll_id = conn.last_insert_rowid();
838
839 // Collection members
840 conn.execute(
841 "INSERT INTO collection_members (collection_id, sample_hash, added_at) VALUES (?1, 's1', ?2)",
842 rusqlite::params![coll_id, now],
843 ).unwrap();
844
845 set_sync_state(conn, "applying_remote", "0").unwrap();
846 clear_changelog(conn);
847
848 let total = create_initial_snapshot(conn).unwrap();
849 // 2 samples + 2 tags + 1 collection + 1 collection_member = 6
850 assert_eq!(total, 6);
851
852 assert_eq!(changelog_count(conn, Some("samples"), Some("INSERT")), 2);
853 assert_eq!(changelog_count(conn, Some("tags"), Some("INSERT")), 2);
854 assert_eq!(changelog_count(conn, Some("collections"), Some("INSERT")), 1);
855 assert_eq!(changelog_count(conn, Some("collection_members"), Some("INSERT")), 1);
856 }
857
858 #[test]
859 fn snapshot_changelog_entries_contain_valid_json() {
860 let db = setup_test_db();
861 let conn = db.conn();
862
863 set_sync_state(conn, "applying_remote", "1").unwrap();
864 insert_sample(conn, "json_test", "pad.wav", "wav");
865 set_sync_state(conn, "applying_remote", "0").unwrap();
866 clear_changelog(conn);
867
868 create_initial_snapshot(conn).unwrap();
869
870 let data: String = conn.query_row(
871 "SELECT data FROM sync_changelog WHERE table_name = 'samples' AND row_id = 'json_test'",
872 [],
873 |row| row.get(0),
874 ).unwrap();
875
876 let parsed: serde_json::Value = serde_json::from_str(&data).unwrap();
877 assert_eq!(parsed["hash"], "json_test");
878 assert_eq!(parsed["original_name"], "pad.wav");
879 assert_eq!(parsed["file_extension"], "wav");
880 assert!(parsed.get("file_size").is_some());
881 assert!(parsed.get("import_date").is_some());
882 // Columns added in migrations 008-009 must be present in snapshot
883 assert!(parsed.get("cloud_only").is_some(), "snapshot missing cloud_only");
884 assert!(parsed.get("duration").is_some(), "snapshot missing duration");
885 }
886
887 #[test]
888 fn snapshot_samples_columns_match_whitelist() {
889 let db = setup_test_db();
890 let conn = db.conn();
891
892 set_sync_state(conn, "applying_remote", "1").unwrap();
893 insert_sample(conn, "col_test", "check.wav", "wav");
894 set_sync_state(conn, "applying_remote", "0").unwrap();
895 clear_changelog(conn);
896
897 create_initial_snapshot(conn).unwrap();
898
899 let data: String = conn.query_row(
900 "SELECT data FROM sync_changelog WHERE table_name = 'samples' AND row_id = 'col_test'",
901 [],
902 |row| row.get(0),
903 ).unwrap();
904 let parsed: serde_json::Value = serde_json::from_str(&data).unwrap();
905 let obj = parsed.as_object().unwrap();
906
907 // Every column in the whitelist must appear in the snapshot JSON
908 for col in table_columns("samples").unwrap() {
909 assert!(obj.contains_key(*col), "snapshot samples missing column: {}", col);
910 }
911 }
912
913 #[test]
914 fn snapshot_audio_analysis_columns_match_whitelist() {
915 let db = setup_test_db();
916 let conn = db.conn();
917
918 set_sync_state(conn, "applying_remote", "1").unwrap();
919 insert_sample(conn, "aa_test", "tone.wav", "wav");
920 conn.execute(
921 "INSERT INTO audio_analysis (hash, duration, sample_rate, channels, analyzed_at) VALUES ('aa_test', 1.5, 44100, 2, 1000000)",
922 [],
923 ).unwrap();
924 set_sync_state(conn, "applying_remote", "0").unwrap();
925 clear_changelog(conn);
926
927 create_initial_snapshot(conn).unwrap();
928
929 let data: String = conn.query_row(
930 "SELECT data FROM sync_changelog WHERE table_name = 'audio_analysis' AND row_id = 'aa_test'",
931 [],
932 |row| row.get(0),
933 ).unwrap();
934 let parsed: serde_json::Value = serde_json::from_str(&data).unwrap();
935 let obj = parsed.as_object().unwrap();
936
937 for col in table_columns("audio_analysis").unwrap() {
938 assert!(obj.contains_key(*col), "snapshot audio_analysis missing column: {}", col);
939 }
940 }
941
942 #[test]
943 fn loose_files_excluded_from_sync() {
944 let db = setup_test_db();
945 let conn = db.conn();
946 clear_changelog(conn);
947
948 // Insert loose_files — should NOT fire trigger
949 conn.execute(
950 "INSERT INTO user_config (key, value) VALUES ('loose_files', '1')",
951 [],
952 ).unwrap();
953 assert_eq!(changelog_count(conn, Some("user_config"), None), 0);
954
955 // Update loose_files — should NOT fire trigger
956 conn.execute(
957 "UPDATE user_config SET value = '0' WHERE key = 'loose_files'",
958 [],
959 ).unwrap();
960 assert_eq!(changelog_count(conn, Some("user_config"), None), 0);
961
962 // Normal key — should fire trigger
963 conn.execute(
964 "INSERT INTO user_config (key, value) VALUES ('theme', 'dark')",
965 [],
966 ).unwrap();
967 assert_eq!(changelog_count(conn, Some("user_config"), None), 1);
968 }
969
970 #[test]
971 fn snapshot_after_adding_data_does_not_duplicate() {
972 let db = setup_test_db();
973 let conn = db.conn();
974
975 set_sync_state(conn, "applying_remote", "1").unwrap();
976 insert_sample(conn, "first", "one.wav", "wav");
977 set_sync_state(conn, "applying_remote", "0").unwrap();
978 clear_changelog(conn);
979
980 let first = create_initial_snapshot(conn).unwrap();
981 assert_eq!(first, 1);
982
983 // Add more data after snapshot (these go through normal triggers)
984 insert_sample(conn, "second", "two.wav", "wav");
985
986 // Second snapshot call should return 0 (flag already set)
987 let second = create_initial_snapshot(conn).unwrap();
988 assert_eq!(second, 0);
989
990 // Total changelog: 1 from snapshot + 1 from trigger
991 assert_eq!(changelog_count(conn, Some("samples"), None), 2);
992 }
993
994 #[test]
995 fn cleanup_preserves_recent_pushed_entries() {
996 let db = setup_test_db();
997 let conn = db.conn();
998 clear_changelog(conn);
999
1000 // Entry pushed exactly 6 days ago (within 7-day window)
1001 conn.execute(
1002 "INSERT INTO sync_changelog (table_name, op, row_id, timestamp, pushed) \
1003 VALUES ('samples', 'INSERT', 'six_days', datetime('now', '-6 days'), 1)",
1004 [],
1005 ).unwrap();
1006
1007 // Entry pushed exactly 8 days ago (outside 7-day window)
1008 conn.execute(
1009 "INSERT INTO sync_changelog (table_name, op, row_id, timestamp, pushed) \
1010 VALUES ('samples', 'INSERT', 'eight_days', datetime('now', '-8 days'), 1)",
1011 [],
1012 ).unwrap();
1013
1014 let deleted = cleanup_changelog(conn).unwrap();
1015 assert_eq!(deleted, 1);
1016
1017 // six_days should remain, eight_days should be deleted
1018 let remaining_id: String = conn.query_row(
1019 "SELECT row_id FROM sync_changelog",
1020 [],
1021 |r| r.get(0),
1022 ).unwrap();
1023 assert_eq!(remaining_id, "six_days");
1024 }
1025
1026 #[test]
1027 fn enforce_retention_at_exact_cap_is_noop() {
1028 let db = setup_test_db();
1029 let conn = db.conn();
1030 clear_changelog(conn);
1031
1032 for i in 0..MAX_CHANGELOG_ENTRIES {
1033 conn.execute(
1034 "INSERT INTO sync_changelog (table_name, op, row_id, data) \
1035 VALUES ('samples', 'UPDATE', ?1, '{}')",
1036 [format!("row-{}", i)],
1037 ).unwrap();
1038 }
1039
1040 let deleted = enforce_changelog_retention(conn).unwrap();
1041 assert_eq!(deleted, 0);
1042
1043 let count: i64 = conn.query_row(
1044 "SELECT COUNT(*) FROM sync_changelog", [], |r| r.get(0),
1045 ).unwrap();
1046 assert_eq!(count, MAX_CHANGELOG_ENTRIES);
1047 }
1048
1049 #[test]
1050 fn cleanup_and_retention_combined() {
1051 let db = setup_test_db();
1052 let conn = db.conn();
1053 clear_changelog(conn);
1054
1055 // Insert many old pushed entries (should be cleaned by cleanup)
1056 for i in 0..100 {
1057 conn.execute(
1058 "INSERT INTO sync_changelog (table_name, op, row_id, timestamp, pushed) \
1059 VALUES ('samples', 'UPDATE', ?1, datetime('now', '-14 days'), 1)",
1060 [format!("old-{}", i)],
1061 ).unwrap();
1062 }
1063
1064 // Insert entries at the retention cap to test combined behavior
1065 for i in 0..MAX_CHANGELOG_ENTRIES {
1066 conn.execute(
1067 "INSERT INTO sync_changelog (table_name, op, row_id, data) \
1068 VALUES ('samples', 'INSERT', ?1, '{}')",
1069 [format!("new-{}", i)],
1070 ).unwrap();
1071 }
1072
1073 // cleanup removes old pushed entries
1074 let cleaned = cleanup_changelog(conn).unwrap();
1075 assert_eq!(cleaned, 100);
1076
1077 // retention should be a noop now (exactly MAX_CHANGELOG_ENTRIES remain)
1078 let retained = enforce_changelog_retention(conn).unwrap();
1079 assert_eq!(retained, 0);
1080 }
1081
1082 #[test]
1083 fn sync_state_get_and_set_roundtrip() {
1084 let db = setup_test_db();
1085 let conn = db.conn();
1086
1087 set_sync_state(conn, "auto_sync_enabled", "1").unwrap();
1088 assert_eq!(get_sync_state(conn, "auto_sync_enabled").unwrap(), "1");
1089
1090 set_sync_state(conn, "auto_sync_enabled", "0").unwrap();
1091 assert_eq!(get_sync_state(conn, "auto_sync_enabled").unwrap(), "0");
1092 }
1093
1094 // ── Download query logic ──
1095
1096 #[test]
1097 fn missing_blobs_query_finds_sync_enabled_only() {
1098 let db = setup_test_db();
1099 let conn = db.conn();
1100
1101 // Create two VFS: one with sync_files=true, one without
1102 let vfs_sync = insert_vfs(conn, "Synced", true);
1103 let vfs_local = insert_vfs(conn, "Local", false);
1104
1105 insert_sample(conn, "hash_a", "kick.wav", "wav");
1106 insert_sample(conn, "hash_b", "snare.wav", "wav");
1107
1108 // Link hash_a to synced VFS, hash_b to local-only VFS
1109 let now = chrono::Utc::now().timestamp();
1110 conn.execute(
1111 "INSERT INTO vfs_nodes (vfs_id, parent_id, name, node_type, sample_hash, created_at) VALUES (?1, NULL, 'kick.wav', 'sample', 'hash_a', ?2)",
1112 rusqlite::params![vfs_sync, now],
1113 ).unwrap();
1114 conn.execute(
1115 "INSERT INTO vfs_nodes (vfs_id, parent_id, name, node_type, sample_hash, created_at) VALUES (?1, NULL, 'snare.wav', 'sample', 'hash_b', ?2)",
1116 rusqlite::params![vfs_local, now],
1117 ).unwrap();
1118
1119 // Run the same query as download_missing_blobs
1120 let mut stmt = conn.prepare(
1121 "SELECT DISTINCT s.hash, s.file_extension
1122 FROM samples s
1123 JOIN vfs_nodes vn ON vn.sample_hash = s.hash
1124 JOIN vfs v ON v.id = vn.vfs_id
1125 WHERE v.sync_files = 1",
1126 ).unwrap();
1127 let rows: Vec<(String, String)> = stmt
1128 .query_map([], |row| Ok((row.get(0)?, row.get(1)?)))
1129 .unwrap()
1130 .collect::<std::result::Result<Vec<_>, _>>()
1131 .unwrap();
1132
1133 assert_eq!(rows.len(), 1);
1134 assert_eq!(rows[0].0, "hash_a");
1135 }
1136
1137 #[test]
1138 fn missing_blobs_query_deduplicates_multi_vfs() {
1139 let db = setup_test_db();
1140 let conn = db.conn();
1141
1142 let vfs1 = insert_vfs(conn, "Lib1", true);
1143 let vfs2 = insert_vfs(conn, "Lib2", true);
1144 insert_sample(conn, "shared_hash", "shared.wav", "wav");
1145
1146 // Same sample linked in two synced VFS entries
1147 let now = chrono::Utc::now().timestamp();
1148 conn.execute(
1149 "INSERT INTO vfs_nodes (vfs_id, parent_id, name, node_type, sample_hash, created_at) VALUES (?1, NULL, 'a.wav', 'sample', 'shared_hash', ?2)",
1150 rusqlite::params![vfs1, now],
1151 ).unwrap();
1152 conn.execute(
1153 "INSERT INTO vfs_nodes (vfs_id, parent_id, name, node_type, sample_hash, created_at) VALUES (?1, NULL, 'b.wav', 'sample', 'shared_hash', ?2)",
1154 rusqlite::params![vfs2, now],
1155 ).unwrap();
1156
1157 let mut stmt = conn.prepare(
1158 "SELECT DISTINCT s.hash, s.file_extension
1159 FROM samples s
1160 JOIN vfs_nodes vn ON vn.sample_hash = s.hash
1161 JOIN vfs v ON v.id = vn.vfs_id
1162 WHERE v.sync_files = 1",
1163 ).unwrap();
1164 let rows: Vec<(String, String)> = stmt
1165 .query_map([], |row| Ok((row.get(0)?, row.get(1)?)))
1166 .unwrap()
1167 .collect::<std::result::Result<Vec<_>, _>>()
1168 .unwrap();
1169
1170 // DISTINCT should yield exactly one row
1171 assert_eq!(rows.len(), 1);
1172 }
1173
1174 // ── Upload query logic ──
1175
1176 #[test]
1177 fn upload_pending_query_finds_non_cloud_only() {
1178 let db = setup_test_db();
1179 let conn = db.conn();
1180
1181 let vfs = insert_vfs(conn, "Synced", true);
1182 insert_sample(conn, "local_hash", "kick.wav", "wav");
1183
1184 // Mark one sample as cloud_only
1185 conn.execute("UPDATE samples SET cloud_only = 1 WHERE hash = 'local_hash'", []).unwrap();
1186
1187 insert_sample(conn, "present_hash", "snare.wav", "wav");
1188
1189 let now = chrono::Utc::now().timestamp();
1190 conn.execute(
1191 "INSERT INTO vfs_nodes (vfs_id, parent_id, name, node_type, sample_hash, created_at) VALUES (?1, NULL, 'kick.wav', 'sample', 'local_hash', ?2)",
1192 rusqlite::params![vfs, now],
1193 ).unwrap();
1194 conn.execute(
1195 "INSERT INTO vfs_nodes (vfs_id, parent_id, name, node_type, sample_hash, created_at) VALUES (?1, NULL, 'snare.wav', 'sample', 'present_hash', ?2)",
1196 rusqlite::params![vfs, now],
1197 ).unwrap();
1198
1199 // Same query as upload_pending_blobs
1200 let mut stmt = conn.prepare(
1201 "SELECT DISTINCT s.hash, s.file_extension, s.file_size
1202 FROM samples s
1203 JOIN vfs_nodes vn ON vn.sample_hash = s.hash
1204 JOIN vfs v ON v.id = vn.vfs_id
1205 WHERE v.sync_files = 1 AND s.cloud_only = 0",
1206 ).unwrap();
1207 let rows: Vec<(String, String, i64)> = stmt
1208 .query_map([], |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)))
1209 .unwrap()
1210 .collect::<std::result::Result<Vec<_>, _>>()
1211 .unwrap();
1212
1213 // Only present_hash (cloud_only=0)
1214 assert_eq!(rows.len(), 1);
1215 assert_eq!(rows[0].0, "present_hash");
1216 }
1217
1218 #[test]
1219 fn push_changelog_reads_unpushed_in_order() {
1220 let db = setup_test_db();
1221 let conn = db.conn();
1222 clear_changelog(conn);
1223
1224 // Insert changelog entries with specific order
1225 for i in 0..5 {
1226 conn.execute(
1227 "INSERT INTO sync_changelog (table_name, op, row_id, data, pushed) VALUES ('samples', 'INSERT', ?1, '{}', 0)",
1228 [format!("push-{}", i)],
1229 ).unwrap();
1230 }
1231 // Mark one as already pushed
1232 conn.execute(
1233 "UPDATE sync_changelog SET pushed = 1 WHERE row_id = 'push-2'",
1234 [],
1235 ).unwrap();
1236
1237 // Same query as push_changes
1238 let mut stmt = conn.prepare(
1239 "SELECT id, table_name, op, row_id, timestamp, data
1240 FROM sync_changelog
1241 WHERE pushed = 0
1242 ORDER BY id ASC
1243 LIMIT ?1",
1244 ).unwrap();
1245 let rows: Vec<(i64, String, String, String)> = stmt
1246 .query_map([500i64], |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?)))
1247 .unwrap()
1248 .collect::<std::result::Result<Vec<_>, _>>()
1249 .unwrap();
1250
1251 // Should skip push-2 (already pushed)
1252 assert_eq!(rows.len(), 4);
1253 let row_ids: Vec<&str> = rows.iter().map(|r| r.3.as_str()).collect();
1254 assert!(!row_ids.contains(&"push-2"));
1255 // Should be in ASC order
1256 assert_eq!(row_ids[0], "push-0");
1257 assert_eq!(row_ids[3], "push-4");
1258 }
1259
1260 // ── Resolve: mixed upsert/delete ordering ──
1261
1262 #[test]
1263 fn apply_remote_changes_mixed_ops_correct_order() {
1264 let db = setup_test_db();
1265 let conn = db.conn();
1266 clear_changelog(conn);
1267
1268 // Insert a sample first
1269 let changes_insert = vec![
1270 change("samples", ChangeOp::Insert, "mix_hash", Some(json!({
1271 "hash": "mix_hash", "original_name": "mixed.wav",
1272 "file_extension": "wav", "file_size": 1024,
1273 "import_date": 1000000, "last_modified": 1000000,
1274 "cloud_only": 0
1275 }))),
1276 change("tags", ChangeOp::Insert, "mix_hash:bass", Some(json!({
1277 "sample_hash": "mix_hash", "tag": "bass"
1278 }))),
1279 ];
1280 apply_remote_changes(conn, &changes_insert).unwrap();
1281
1282 // Now delete tag then sample (mixed batch)
1283 let changes_delete = vec![
1284 change("tags", ChangeOp::Delete, "mix_hash:bass", None),
1285 change("samples", ChangeOp::Delete, "mix_hash", None),
1286 ];
1287 let applied = apply_remote_changes(conn, &changes_delete).unwrap();
1288 assert_eq!(applied, 2);
1289
1290 let sample_count: i64 = conn.query_row(
1291 "SELECT COUNT(*) FROM samples WHERE hash = 'mix_hash'",
1292 [], |row| row.get(0),
1293 ).unwrap();
1294 assert_eq!(sample_count, 0);
1295
1296 let tag_count: i64 = conn.query_row(
1297 "SELECT COUNT(*) FROM tags WHERE sample_hash = 'mix_hash'",
1298 [], |row| row.get(0),
1299 ).unwrap();
1300 assert_eq!(tag_count, 0);
1301 }
1302
1303 #[test]
1304 fn apply_upsert_updates_existing_row() {
1305 let db = setup_test_db();
1306 let conn = db.conn();
1307
1308 // Insert initial sample
1309 apply_upsert(conn, "samples", &json!({
1310 "hash": "upd_hash", "original_name": "old.wav",
1311 "file_extension": "wav", "file_size": 1024,
1312 "import_date": 1000000, "last_modified": 1000000,
1313 "cloud_only": 0
1314 })).unwrap();
1315
1316 // Update via upsert (ON CONFLICT DO UPDATE)
1317 apply_upsert(conn, "samples", &json!({
1318 "hash": "upd_hash", "original_name": "new.wav",
1319 "file_extension": "wav", "file_size": 2048,
1320 "import_date": 1000000, "last_modified": 2000000,
1321 "cloud_only": 0
1322 })).unwrap();
1323
1324 let name: String = conn.query_row(
1325 "SELECT original_name FROM samples WHERE hash = 'upd_hash'",
1326 [], |row| row.get(0),
1327 ).unwrap();
1328 assert_eq!(name, "new.wav");
1329
1330 let size: i64 = conn.query_row(
1331 "SELECT file_size FROM samples WHERE hash = 'upd_hash'",
1332 [], |row| row.get(0),
1333 ).unwrap();
1334 assert_eq!(size, 2048);
1335
1336 // Should still be only 1 row
1337 let count: i64 = conn.query_row(
1338 "SELECT COUNT(*) FROM samples WHERE hash = 'upd_hash'",
1339 [], |row| row.get(0),
1340 ).unwrap();
1341 assert_eq!(count, 1);
1342 }
1343
1344 #[test]
1345 fn apply_delete_nonexistent_is_noop() {
1346 let db = setup_test_db();
1347 let conn = db.conn();
1348
1349 // Delete a hash that doesn't exist — should succeed (no-op)
1350 let result = apply_delete(conn, "samples", "nonexistent_hash");
1351 assert!(result.is_ok());
1352 }
1353 }
1354