Skip to main content

max / audiofiles

54.3 KB · 1808 lines History Blame Raw
1 //! Tests for [`super`].
2
3 use super::super::{
4 DELETE_ORDER, UPSERT_ORDER, pk_columns,
5 resolve::{apply_delete, apply_remote_changes, apply_upsert},
6 table_columns,
7 };
8 use super::*;
9 use audiofiles_core::db::Database;
10 use serde_json::json;
11 use synckit_client::{ChangeEntry, ChangeOp};
12
13 fn setup_test_db() -> Database {
14 Database::open_in_memory().expect("Failed to create test DB")
15 }
16
17 fn insert_sample(conn: &Connection, hash: &str, name: &str, ext: &str) {
18 let now = chrono::Utc::now().timestamp();
19 conn.execute(
20 "INSERT INTO samples (hash, original_name, file_extension, file_size, import_date, last_modified) VALUES (?1, ?2, ?3, 1024, ?4, ?4)",
21 rusqlite::params![hash, name, ext, now],
22 ).unwrap();
23 }
24
25 fn insert_vfs(conn: &Connection, name: &str, sync_files: bool) -> i64 {
26 let now = chrono::Utc::now().timestamp();
27 conn.execute(
28 "INSERT INTO vfs (name, created_at, modified_at, sync_files) VALUES (?1, ?2, ?2, ?3)",
29 rusqlite::params![name, now, sync_files as i64],
30 )
31 .unwrap();
32 conn.last_insert_rowid()
33 }
34
35 fn clear_changelog(conn: &Connection) {
36 conn.execute("DELETE FROM sync_changelog", []).unwrap();
37 }
38
39 fn changelog_count(conn: &Connection, table: Option<&str>, op: Option<&str>) -> i64 {
40 match (table, op) {
41 (Some(t), Some(o)) => conn
42 .query_row(
43 "SELECT COUNT(*) FROM sync_changelog WHERE table_name = ?1 AND op = ?2",
44 rusqlite::params![t, o],
45 |row| row.get(0),
46 )
47 .unwrap(),
48 (Some(t), None) => conn
49 .query_row(
50 "SELECT COUNT(*) FROM sync_changelog WHERE table_name = ?1",
51 [t],
52 |row| row.get(0),
53 )
54 .unwrap(),
55 (None, Some(o)) => conn
56 .query_row(
57 "SELECT COUNT(*) FROM sync_changelog WHERE op = ?1",
58 [o],
59 |row| row.get(0),
60 )
61 .unwrap(),
62 (None, None) => conn
63 .query_row("SELECT COUNT(*) FROM sync_changelog", [], |row| row.get(0))
64 .unwrap(),
65 }
66 }
67
68 fn change(table: &str, op: ChangeOp, row_id: &str, data: Option<serde_json::Value>) -> ChangeEntry {
69 ChangeEntry {
70 table: table.to_string(),
71 op,
72 row_id: row_id.to_string(),
73 timestamp: chrono::Utc::now(),
74 hlc: synckit_client::Hlc::zero(synckit_client::DeviceId::nil()),
75 data,
76 extra: serde_json::Map::default(),
77 }
78 }
79
80 // FK ordering
81
82 #[test]
83 fn upsert_order_parents_before_children() {
84 let pos = |t: &str| UPSERT_ORDER.iter().position(|x| *x == t).unwrap();
85 assert!(pos("vfs") < pos("vfs_nodes"));
86 assert!(pos("samples") < pos("audio_analysis"));
87 assert!(pos("samples") < pos("tags"));
88 assert!(pos("samples") < pos("collection_members"));
89 assert!(pos("collections") < pos("collection_members"));
90 }
91
92 #[test]
93 fn delete_order_children_before_parents() {
94 let pos = |t: &str| DELETE_ORDER.iter().position(|x| *x == t).unwrap();
95 assert!(pos("vfs_nodes") < pos("vfs"));
96 assert!(pos("audio_analysis") < pos("samples"));
97 assert!(pos("tags") < pos("samples"));
98 assert!(pos("collection_members") < pos("collections"));
99 assert!(pos("collection_members") < pos("samples"));
100 }
101
102 #[test]
103 fn orders_are_exact_reverses() {
104 let reversed: Vec<&str> = UPSERT_ORDER.iter().rev().copied().collect();
105 assert_eq!(reversed, DELETE_ORDER);
106 }
107
108 // Column whitelists
109
110 #[test]
111 fn all_tables_have_column_whitelists() {
112 for table in UPSERT_ORDER {
113 assert!(
114 table_columns(table).is_some(),
115 "missing column whitelist for: {table}"
116 );
117 }
118 }
119
120 #[test]
121 fn unknown_table_returns_none() {
122 assert!(table_columns("nonexistent").is_none());
123 assert!(table_columns("fingerprints").is_none());
124 }
125
126 #[test]
127 fn pk_columns_covers_all_tables() {
128 for table in UPSERT_ORDER {
129 let pks = pk_columns(table);
130 assert!(!pks.is_empty(), "missing pk_columns for: {table}");
131 }
132 assert_eq!(pk_columns("tags"), &["sample_hash", "tag"]);
133 assert_eq!(
134 pk_columns("collection_members"),
135 &["collection_id", "sample_hash"]
136 );
137 }
138
139 // Triggers
140
141 #[test]
142 fn sample_insert_fires_trigger() {
143 let db = setup_test_db();
144 let conn = db.conn();
145 clear_changelog(conn);
146
147 insert_sample(conn, "abc123", "kick.wav", "wav");
148
149 assert_eq!(changelog_count(conn, Some("samples"), Some("INSERT")), 1);
150
151 // Post-M018: row_id is `hash_row_id(salt, "abc123")` so we look up
152 // by table+op and inspect the canonical hash inside the encrypted-
153 // at-the-wire `data` field.
154 let data: String = conn
155 .query_row(
156 "SELECT data FROM sync_changelog WHERE table_name = 'samples' AND op = 'INSERT'",
157 [],
158 |row| row.get(0),
159 )
160 .unwrap();
161 let parsed: serde_json::Value = serde_json::from_str(&data).unwrap();
162 assert!(parsed.get("cloud_only").is_some());
163 assert_eq!(parsed["hash"], "abc123");
164 }
165
166 #[test]
167 fn vfs_insert_fires_trigger() {
168 let db = setup_test_db();
169 let conn = db.conn();
170 clear_changelog(conn);
171
172 let vfs_id = insert_vfs(conn, "Library", true);
173
174 assert_eq!(changelog_count(conn, Some("vfs"), Some("INSERT")), 1);
175
176 let row_id: String = conn
177 .query_row(
178 "SELECT row_id FROM sync_changelog WHERE table_name = 'vfs'",
179 [],
180 |row| row.get(0),
181 )
182 .unwrap();
183 assert_eq!(row_id, vfs_id.to_string());
184 }
185
186 #[test]
187 fn tag_insert_fires_trigger() {
188 let db = setup_test_db();
189 let conn = db.conn();
190 insert_sample(conn, "hash1", "snare.wav", "wav");
191 clear_changelog(conn);
192
193 conn.execute(
194 "INSERT INTO tags (sample_hash, tag) VALUES ('hash1', 'drums')",
195 [],
196 )
197 .unwrap();
198
199 assert_eq!(changelog_count(conn, Some("tags"), Some("INSERT")), 1);
200
201 // Post-M018: row_id is hashed; the cleartext key lives in `data`.
202 let data: String = conn
203 .query_row(
204 "SELECT data FROM sync_changelog WHERE table_name = 'tags'",
205 [],
206 |row| row.get(0),
207 )
208 .unwrap();
209 let parsed: serde_json::Value = serde_json::from_str(&data).unwrap();
210 assert_eq!(parsed["sample_hash"], "hash1");
211 assert_eq!(parsed["tag"], "drums");
212 }
213
214 #[test]
215 fn sample_features_insert_fires_trigger() {
216 let db = setup_test_db();
217 let conn = db.conn();
218 insert_sample(conn, "hf", "kick.wav", "wav");
219 clear_changelog(conn);
220
221 conn.execute(
222 "INSERT INTO sample_features (hash, feat_version, vector, computed_at) \
223 VALUES ('hf', 1, '[1.0,2.0]', 100)",
224 [],
225 )
226 .unwrap();
227
228 assert_eq!(
229 changelog_count(conn, Some("sample_features"), Some("INSERT")),
230 1
231 );
232
233 let data: String = conn
234 .query_row(
235 "SELECT data FROM sync_changelog WHERE table_name = 'sample_features'",
236 [],
237 |row| row.get(0),
238 )
239 .unwrap();
240 let parsed: serde_json::Value = serde_json::from_str(&data).unwrap();
241 assert_eq!(parsed["hash"], "hf");
242 assert_eq!(parsed["feat_version"], 1);
243 // The JSON-array vector embeds as a string payload.
244 assert_eq!(parsed["vector"], "[1.0,2.0]");
245 }
246
247 #[test]
248 fn tag_rules_insert_fires_trigger() {
249 let db = setup_test_db();
250 let conn = db.conn();
251 clear_changelog(conn);
252
253 conn.execute(
254 "INSERT INTO tag_rules (id, name, enabled, priority, match_mode, conditions, actions, created_at) \
255 VALUES ('r1', 'kicks', 1, 0, '\"all\"', '[]', '[]', 0)",
256 [],
257 ).unwrap();
258
259 assert_eq!(changelog_count(conn, Some("tag_rules"), Some("INSERT")), 1);
260 let row_id: String = conn
261 .query_row(
262 "SELECT row_id FROM sync_changelog WHERE table_name = 'tag_rules'",
263 [],
264 |row| row.get(0),
265 )
266 .unwrap();
267 // Opaque rule id is non-sensitive -> cleartext row_id.
268 assert_eq!(row_id, "r1");
269 }
270
271 #[test]
272 fn tag_provenance_insert_fires_trigger() {
273 let db = setup_test_db();
274 let conn = db.conn();
275 insert_sample(conn, "ph", "kick.wav", "wav");
276 clear_changelog(conn);
277
278 conn.execute(
279 "INSERT INTO tag_provenance (sample_hash, tag, source, rule_id) \
280 VALUES ('ph', 'instrument.drum.kick', 'rule', 'r1')",
281 [],
282 )
283 .unwrap();
284
285 assert_eq!(
286 changelog_count(conn, Some("tag_provenance"), Some("INSERT")),
287 1
288 );
289 let data: String = conn
290 .query_row(
291 "SELECT data FROM sync_changelog WHERE table_name = 'tag_provenance'",
292 [],
293 |row| row.get(0),
294 )
295 .unwrap();
296 let parsed: serde_json::Value = serde_json::from_str(&data).unwrap();
297 assert_eq!(parsed["sample_hash"], "ph");
298 assert_eq!(parsed["tag"], "instrument.drum.kick");
299 assert_eq!(parsed["source"], "rule");
300 assert_eq!(parsed["rule_id"], "r1");
301 }
302
303 #[test]
304 fn collection_member_insert_fires_trigger() {
305 let db = setup_test_db();
306 let conn = db.conn();
307 insert_sample(conn, "hash2", "hat.wav", "wav");
308 let now = chrono::Utc::now().timestamp();
309 conn.execute(
310 "INSERT INTO collections (name, description, created_at) VALUES ('Faves', NULL, ?1)",
311 [now],
312 )
313 .unwrap();
314 let collection_id = conn.last_insert_rowid();
315 clear_changelog(conn);
316
317 conn.execute(
318 "INSERT INTO collection_members (collection_id, sample_hash, added_at) VALUES (?1, 'hash2', ?2)",
319 rusqlite::params![collection_id, now],
320 ).unwrap();
321
322 assert_eq!(
323 changelog_count(conn, Some("collection_members"), Some("INSERT")),
324 1
325 );
326
327 let data: String = conn
328 .query_row(
329 "SELECT data FROM sync_changelog WHERE table_name = 'collection_members'",
330 [],
331 |row| row.get(0),
332 )
333 .unwrap();
334 let parsed: serde_json::Value = serde_json::from_str(&data).unwrap();
335 assert_eq!(parsed["collection_id"].as_i64().unwrap(), collection_id);
336 assert_eq!(parsed["sample_hash"], "hash2");
337 }
338
339 // Trigger suppression
340
341 #[test]
342 fn trigger_suppression_during_remote_apply() {
343 let db = setup_test_db();
344 let conn = db.conn();
345 clear_changelog(conn);
346
347 set_sync_state(conn, "applying_remote", "1").unwrap();
348 insert_sample(conn, "suppressed", "test.wav", "wav");
349 assert_eq!(changelog_count(conn, None, None), 0);
350
351 set_sync_state(conn, "applying_remote", "0").unwrap();
352 insert_sample(conn, "unsuppressed", "test2.wav", "wav");
353 assert_eq!(changelog_count(conn, None, None), 1);
354 }
355
356 // apply_upsert
357
358 #[test]
359 fn apply_upsert_inserts_sample() {
360 let db = setup_test_db();
361 let conn = db.conn();
362
363 let data = json!({
364 "hash": "upsert_hash",
365 "original_name": "synced.wav",
366 "file_extension": "wav",
367 "file_size": 2048,
368 "import_date": 1_000_000,
369 "last_modified": 1_000_000,
370 "cloud_only": 0
371 });
372
373 apply_upsert(conn, "samples", &data).unwrap();
374
375 let name: String = conn
376 .query_row(
377 "SELECT original_name FROM samples WHERE hash = 'upsert_hash'",
378 [],
379 |row| row.get(0),
380 )
381 .unwrap();
382 assert_eq!(name, "synced.wav");
383 }
384
385 #[test]
386 fn apply_upsert_inserts_vfs_node() {
387 let db = setup_test_db();
388 let conn = db.conn();
389
390 let vfs_id = insert_vfs(conn, "TestVFS", false);
391 insert_sample(conn, "node_hash", "pad.wav", "wav");
392
393 let data = json!({
394 "id": 999,
395 "vfs_id": vfs_id,
396 "parent_id": null,
397 "name": "pad.wav",
398 "node_type": "sample",
399 "sample_hash": "node_hash",
400 "created_at": 1_000_000
401 });
402
403 apply_upsert(conn, "vfs_nodes", &data).unwrap();
404
405 let name: String = conn
406 .query_row("SELECT name FROM vfs_nodes WHERE id = 999", [], |row| {
407 row.get(0)
408 })
409 .unwrap();
410 assert_eq!(name, "pad.wav");
411 }
412
413 #[test]
414 fn apply_upsert_unknown_table_is_no_op() {
415 let db = setup_test_db();
416 let conn = db.conn();
417
418 let data = json!({"id": "abc"});
419 let result = apply_upsert(conn, "nonexistent_table", &data);
420 assert!(result.is_ok());
421 }
422
423 // apply_delete
424
425 #[test]
426 fn apply_delete_samples_tombstones_not_hard_deletes() {
427 let db = setup_test_db();
428 let conn = db.conn();
429 insert_sample(conn, "del_hash", "delete_me.wav", "wav");
430
431 // A remote samples delete must soft-delete (tombstone), never hard-delete:
432 // the row stays so the engine-level CASCADE never fires.
433 apply_delete(conn, "samples", "del_hash", None).unwrap();
434
435 let (count, tombstoned): (i64, i64) = conn
436 .query_row(
437 "SELECT COUNT(*), COUNT(deleted_at) FROM samples WHERE hash = 'del_hash'",
438 [],
439 |row| Ok((row.get(0)?, row.get(1)?)),
440 )
441 .unwrap();
442 assert_eq!(count, 1, "row must survive a remote delete");
443 assert_eq!(tombstoned, 1, "row must be tombstoned (deleted_at set)");
444 }
445
446 #[test]
447 fn apply_delete_samples_preserves_placements_and_tags() {
448 let db = setup_test_db();
449 let conn = db.conn();
450 let vfs_id = insert_vfs(conn, "Library", true);
451 insert_sample(conn, "keep_hash", "kick.wav", "wav");
452 conn.execute(
453 "INSERT INTO vfs_nodes (vfs_id, parent_id, name, node_type, sample_hash, created_at) \
454 VALUES (?1, NULL, 'kick.wav', 'sample', 'keep_hash', 1000)",
455 [vfs_id],
456 )
457 .unwrap();
458 conn.execute(
459 "INSERT INTO tags (sample_hash, tag) VALUES ('keep_hash', 'drums')",
460 [],
461 )
462 .unwrap();
463
464 // The M018+ wire form: hash lives in `data`, row_id is opaque.
465 apply_delete(
466 conn,
467 "samples",
468 "opaque_row_id",
469 Some(&json!({ "hash": "keep_hash" })),
470 )
471 .unwrap();
472
473 let placements: i64 = conn
474 .query_row(
475 "SELECT COUNT(*) FROM vfs_nodes WHERE sample_hash = 'keep_hash'",
476 [],
477 |r| r.get(0),
478 )
479 .unwrap();
480 let tags: i64 = conn
481 .query_row(
482 "SELECT COUNT(*) FROM tags WHERE sample_hash = 'keep_hash'",
483 [],
484 |r| r.get(0),
485 )
486 .unwrap();
487 assert_eq!(
488 placements, 1,
489 "remote delete must not cascade-wipe placements"
490 );
491 assert_eq!(tags, 1, "remote delete must not cascade-wipe tags");
492 }
493
494 #[test]
495 fn apply_delete_composite_pk_tags() {
496 let db = setup_test_db();
497 let conn = db.conn();
498 insert_sample(conn, "tag_hash", "tagged.wav", "wav");
499 conn.execute(
500 "INSERT INTO tags (sample_hash, tag) VALUES ('tag_hash', 'bass')",
501 [],
502 )
503 .unwrap();
504
505 apply_delete(conn, "tags", "tag_hash:bass", None).unwrap();
506
507 let count: i64 = conn
508 .query_row(
509 "SELECT COUNT(*) FROM tags WHERE sample_hash = 'tag_hash' AND tag = 'bass'",
510 [],
511 |row| row.get(0),
512 )
513 .unwrap();
514 assert_eq!(count, 0);
515 }
516
517 // Full pipeline
518
519 #[test]
520 fn apply_remote_changes_full_pipeline() {
521 let db = setup_test_db();
522 let conn = db.conn();
523 clear_changelog(conn);
524
525 // Changes in wrong FK order, apply_remote_changes reorders via UPSERT_ORDER
526 let changes = vec![
527 change(
528 "vfs_nodes",
529 ChangeOp::Insert,
530 "1",
531 Some(json!({
532 "id": 1, "vfs_id": 1, "parent_id": null,
533 "name": "kick.wav", "node_type": "sample",
534 "sample_hash": "pipe_hash", "created_at": 1_000_000
535 })),
536 ),
537 change(
538 "vfs",
539 ChangeOp::Insert,
540 "1",
541 Some(json!({
542 "id": 1, "name": "Library", "created_at": 1_000_000,
543 "modified_at": 1_000_000, "sync_files": 1
544 })),
545 ),
546 change(
547 "samples",
548 ChangeOp::Insert,
549 "pipe_hash",
550 Some(json!({
551 "hash": "pipe_hash", "original_name": "kick.wav",
552 "file_extension": "wav", "file_size": 4096,
553 "import_date": 1_000_000, "last_modified": 1_000_000,
554 "cloud_only": 0
555 })),
556 ),
557 ];
558
559 let applied = apply_remote_changes(conn, &changes).unwrap();
560 assert_eq!(applied, 3);
561
562 let sample_count: i64 = conn
563 .query_row(
564 "SELECT COUNT(*) FROM samples WHERE hash = 'pipe_hash'",
565 [],
566 |row| row.get(0),
567 )
568 .unwrap();
569 assert_eq!(sample_count, 1);
570
571 let vfs_count: i64 = conn
572 .query_row(
573 "SELECT COUNT(*) FROM vfs WHERE name = 'Library'",
574 [],
575 |row| row.get(0),
576 )
577 .unwrap();
578 assert_eq!(vfs_count, 1);
579
580 let node_count: i64 = conn
581 .query_row(
582 "SELECT COUNT(*) FROM vfs_nodes WHERE sample_hash = 'pipe_hash'",
583 [],
584 |row| row.get(0),
585 )
586 .unwrap();
587 assert_eq!(node_count, 1);
588
589 // Changelog should be empty (triggers were suppressed)
590 assert_eq!(changelog_count(conn, None, None), 0);
591 }
592
593 // Initial snapshot
594
595 #[test]
596 fn create_initial_snapshot_captures_all_rows() {
597 let db = setup_test_db();
598 let conn = db.conn();
599
600 // Suppress triggers during data setup
601 set_sync_state(conn, "applying_remote", "1").unwrap();
602 insert_sample(conn, "snap1", "one.wav", "wav");
603 insert_sample(conn, "snap2", "two.wav", "wav");
604 insert_vfs(conn, "TestLib", false);
605 set_sync_state(conn, "applying_remote", "0").unwrap();
606
607 clear_changelog(conn);
608
609 let total = create_initial_snapshot(conn).unwrap();
610 // 2 samples + 1 vfs + 1 user_config seed
611 // (`sample_tombstone_retain_days = 30`, seeded by M019).
612 assert_eq!(total, 4);
613 }
614
615 #[test]
616 fn create_initial_snapshot_excludes_device_local_config_keys() {
617 // Third export path (fuzz-2026-07-21 #3): the snapshot writes changelog
618 // rows directly, so it must apply the registry-generated filter too. A
619 // device-local key set locally must never appear in the snapshot; a
620 // replicated one must.
621 let db = setup_test_db();
622 let conn = db.conn();
623
624 set_sync_state(conn, "applying_remote", "1").unwrap();
625 for key in ["mirror_path", "mirror_enabled", "import_preflight_disabled"] {
626 conn.execute(
627 "INSERT OR REPLACE INTO user_config (key, value) VALUES (?1, '1')",
628 [key],
629 )
630 .unwrap();
631 }
632 conn.execute(
633 "INSERT OR REPLACE INTO user_config (key, value) VALUES ('theme', 'dark')",
634 [],
635 )
636 .unwrap();
637 set_sync_state(conn, "applying_remote", "0").unwrap();
638 clear_changelog(conn);
639
640 create_initial_snapshot(conn).unwrap();
641
642 let snapshotted = |key: &str| -> i64 {
643 conn.query_row(
644 "SELECT COUNT(*) FROM sync_changelog WHERE table_name = 'user_config' AND row_id = ?1",
645 [key],
646 |r| r.get(0),
647 )
648 .unwrap()
649 };
650 assert_eq!(snapshotted("mirror_path"), 0);
651 assert_eq!(snapshotted("mirror_enabled"), 0);
652 assert_eq!(snapshotted("import_preflight_disabled"), 0);
653 assert_eq!(snapshotted("theme"), 1);
654 }
655
656 #[test]
657 fn create_initial_snapshot_idempotent() {
658 let db = setup_test_db();
659 let conn = db.conn();
660
661 set_sync_state(conn, "applying_remote", "1").unwrap();
662 insert_sample(conn, "idem", "test.wav", "wav");
663 set_sync_state(conn, "applying_remote", "0").unwrap();
664 clear_changelog(conn);
665
666 let first = create_initial_snapshot(conn).unwrap();
667 assert!(first > 0);
668
669 let second = create_initial_snapshot(conn).unwrap();
670 assert_eq!(second, 0);
671 }
672
673 // Changelog helpers
674
675 #[test]
676 fn cleanup_removes_old_pushed() {
677 let db = setup_test_db();
678 let conn = db.conn();
679 clear_changelog(conn);
680
681 // Old pushed, should be deleted
682 conn.execute(
683 "INSERT INTO sync_changelog (table_name, op, row_id, timestamp, pushed) \
684 VALUES ('samples', 'INSERT', 'old', datetime('now', '-10 days'), 1)",
685 [],
686 )
687 .unwrap();
688
689 // Recent pushed, should remain
690 conn.execute(
691 "INSERT INTO sync_changelog (table_name, op, row_id, timestamp, pushed) \
692 VALUES ('samples', 'INSERT', 'recent', datetime('now'), 1)",
693 [],
694 )
695 .unwrap();
696
697 // Old unpushed, should remain (never delete unpushed)
698 conn.execute(
699 "INSERT INTO sync_changelog (table_name, op, row_id, timestamp, pushed) \
700 VALUES ('samples', 'INSERT', 'unpushed', datetime('now', '-10 days'), 0)",
701 [],
702 )
703 .unwrap();
704
705 let deleted = cleanup_changelog(conn).unwrap();
706 assert_eq!(deleted, 1);
707
708 let remaining: i64 = conn
709 .query_row("SELECT COUNT(*) FROM sync_changelog", [], |row| row.get(0))
710 .unwrap();
711 assert_eq!(remaining, 2);
712 }
713
714 #[test]
715 fn enforce_retention_caps_entries() {
716 let db = setup_test_db();
717 let conn = db.conn();
718 clear_changelog(conn);
719
720 let total = MAX_CHANGELOG_ENTRIES + 500;
721 for i in 0..total {
722 conn.execute(
723 "INSERT INTO sync_changelog (table_name, op, row_id, data) \
724 VALUES ('samples', 'UPDATE', ?1, '{}')",
725 [format!("row-{i}")],
726 )
727 .unwrap();
728 }
729
730 let before: i64 = conn
731 .query_row("SELECT COUNT(*) FROM sync_changelog", [], |r| r.get(0))
732 .unwrap();
733 assert_eq!(before, total);
734
735 let outcome = enforce_changelog_retention(conn).unwrap();
736 assert_eq!(outcome.deleted, 500);
737 // Every seeded entry was unpushed, so retention had to drop unpushed rows,
738 // the scheduler surfaces this as a hard error.
739 assert_eq!(outcome.unpushed_dropped, 500);
740
741 let after: i64 = conn
742 .query_row("SELECT COUNT(*) FROM sync_changelog", [], |r| r.get(0))
743 .unwrap();
744 assert_eq!(after, MAX_CHANGELOG_ENTRIES);
745
746 // Oldest entries removed, the remaining start at row-500
747 let min_row_id: String = conn
748 .query_row(
749 "SELECT row_id FROM sync_changelog ORDER BY id ASC LIMIT 1",
750 [],
751 |r| r.get(0),
752 )
753 .unwrap();
754 assert_eq!(min_row_id, "row-500");
755 }
756
757 #[test]
758 fn enforce_retention_noop_under_cap() {
759 let db = setup_test_db();
760 let conn = db.conn();
761 clear_changelog(conn);
762
763 for i in 0..5 {
764 conn.execute(
765 "INSERT INTO sync_changelog (table_name, op, row_id, data) \
766 VALUES ('samples', 'INSERT', ?1, '{}')",
767 [format!("row-{i}")],
768 )
769 .unwrap();
770 }
771
772 let outcome = enforce_changelog_retention(conn).unwrap();
773 assert_eq!(outcome.deleted, 0);
774 assert_eq!(outcome.unpushed_dropped, 0);
775 }
776
777 #[test]
778 fn enforce_retention_prefers_pushed_over_unpushed() {
779 let db = setup_test_db();
780 let conn = db.conn();
781 clear_changelog(conn);
782
783 // Seed exactly the cap in already-pushed rows, then 300 unpushed on top.
784 // Retention must remove the 300 excess entirely from the pushed pool and
785 // leave every unpushed row intact, no silent divergence.
786 for i in 0..MAX_CHANGELOG_ENTRIES {
787 conn.execute(
788 "INSERT INTO sync_changelog (table_name, op, row_id, data, pushed) \
789 VALUES ('samples', 'UPDATE', ?1, '{}', 1)",
790 [format!("pushed-{i}")],
791 )
792 .unwrap();
793 }
794 for i in 0..300 {
795 conn.execute(
796 "INSERT INTO sync_changelog (table_name, op, row_id, data, pushed) \
797 VALUES ('samples', 'INSERT', ?1, '{}', 0)",
798 [format!("unpushed-{i}")],
799 )
800 .unwrap();
801 }
802
803 let outcome = enforce_changelog_retention(conn).unwrap();
804 assert_eq!(outcome.deleted, 300);
805 assert_eq!(outcome.unpushed_dropped, 0);
806
807 let unpushed_left: i64 = conn
808 .query_row(
809 "SELECT COUNT(*) FROM sync_changelog WHERE pushed = 0",
810 [],
811 |r| r.get(0),
812 )
813 .unwrap();
814 assert_eq!(unpushed_left, 300);
815 }
816
817 #[test]
818 fn count_pending_counts_unpushed() {
819 let db = setup_test_db();
820 let conn = db.conn();
821 clear_changelog(conn);
822
823 // 3 unpushed
824 for i in 0..3 {
825 conn.execute(
826 "INSERT INTO sync_changelog (table_name, op, row_id, pushed) VALUES ('samples', 'INSERT', ?1, 0)",
827 [format!("un-{i}")],
828 ).unwrap();
829 }
830
831 // 2 pushed
832 for i in 0..2 {
833 conn.execute(
834 "INSERT INTO sync_changelog (table_name, op, row_id, pushed) VALUES ('samples', 'INSERT', ?1, 1)",
835 [format!("push-{i}")],
836 ).unwrap();
837 }
838
839 let count = super::super::count_pending_changes(conn).unwrap();
840 assert_eq!(count, 3);
841 }
842
843 // --- mark_cloud_only_samples tests ---
844
845 #[test]
846 fn mark_cloud_only_marks_missing_blobs() {
847 let db = setup_test_db();
848 let conn = db.conn();
849
850 // Insert two samples, neither has a file on disk
851 insert_sample(conn, "aaa", "kick.wav", "wav");
852 insert_sample(conn, "bbb", "snare.wav", "wav");
853
854 // Use a temp dir as content_dir (empty, no blobs)
855 let tmp = tempfile::tempdir().unwrap();
856
857 let marked = mark_cloud_only_samples(conn, tmp.path()).unwrap();
858 assert_eq!(marked, 2);
859
860 // Verify both are now cloud_only=1
861 let count: i64 = conn
862 .query_row(
863 "SELECT COUNT(*) FROM samples WHERE cloud_only = 1",
864 [],
865 |r| r.get(0),
866 )
867 .unwrap();
868 assert_eq!(count, 2);
869 }
870
871 #[test]
872 fn mark_cloud_only_skips_existing_blobs() {
873 let db = setup_test_db();
874 let conn = db.conn();
875
876 insert_sample(conn, "aaa", "kick.wav", "wav");
877 insert_sample(conn, "bbb", "snare.wav", "wav");
878
879 // Create one blob file, leave other missing
880 let tmp = tempfile::tempdir().unwrap();
881 std::fs::write(tmp.path().join("aaa.wav"), b"fake audio data").unwrap();
882
883 let marked = mark_cloud_only_samples(conn, tmp.path()).unwrap();
884 assert_eq!(marked, 1); // only bbb
885
886 // aaa still cloud_only=0, bbb is cloud_only=1
887 let aaa_co: i32 = conn
888 .query_row(
889 "SELECT cloud_only FROM samples WHERE hash = 'aaa'",
890 [],
891 |r| r.get(0),
892 )
893 .unwrap();
894 assert_eq!(aaa_co, 0);
895
896 let bbb_co: i32 = conn
897 .query_row(
898 "SELECT cloud_only FROM samples WHERE hash = 'bbb'",
899 [],
900 |r| r.get(0),
901 )
902 .unwrap();
903 assert_eq!(bbb_co, 1);
904 }
905
906 #[test]
907 fn mark_cloud_only_suppresses_changelog() {
908 let db = setup_test_db();
909 let conn = db.conn();
910
911 insert_sample(conn, "aaa", "kick.wav", "wav");
912 clear_changelog(conn);
913
914 let tmp = tempfile::tempdir().unwrap();
915 mark_cloud_only_samples(conn, tmp.path()).unwrap();
916
917 // The UPDATE should not appear in changelog (applying_remote suppression)
918 let count = changelog_count(conn, Some("samples"), Some("UPDATE"));
919 assert_eq!(count, 0);
920 }
921
922 #[test]
923 fn mark_cloud_only_skips_loose_files_samples() {
924 let db = setup_test_db();
925 let conn = db.conn();
926
927 // A loose-files sample has source_path set and no blob in content_dir,
928 // it must NOT be flagged cloud_only just because the blob is absent.
929 insert_sample(conn, "aaa", "kick.wav", "wav");
930 conn.execute(
931 "UPDATE samples SET source_path = '/music/kick.wav' WHERE hash = 'aaa'",
932 [],
933 )
934 .unwrap();
935
936 let tmp = tempfile::tempdir().unwrap();
937 let marked = mark_cloud_only_samples(conn, tmp.path()).unwrap();
938 assert_eq!(
939 marked, 0,
940 "loose-files sample must not be marked cloud_only"
941 );
942
943 let aaa_co: i32 = conn
944 .query_row(
945 "SELECT cloud_only FROM samples WHERE hash = 'aaa'",
946 [],
947 |r| r.get(0),
948 )
949 .unwrap();
950 assert_eq!(aaa_co, 0);
951 }
952
953 #[test]
954 fn mark_cloud_only_idempotent() {
955 let db = setup_test_db();
956 let conn = db.conn();
957
958 insert_sample(conn, "aaa", "kick.wav", "wav");
959 let tmp = tempfile::tempdir().unwrap();
960
961 // First call marks it
962 assert_eq!(mark_cloud_only_samples(conn, tmp.path()).unwrap(), 1);
963 // Second call: already cloud_only=1, query only selects cloud_only=0
964 assert_eq!(mark_cloud_only_samples(conn, tmp.path()).unwrap(), 0);
965 }
966
967 // Integration: multi-table snapshot
968
969 #[test]
970 fn snapshot_captures_tags_collections_and_members() {
971 let db = setup_test_db();
972 let conn = db.conn();
973
974 // Suppress triggers during data setup
975 set_sync_state(conn, "applying_remote", "1").unwrap();
976
977 insert_sample(conn, "s1", "kick.wav", "wav");
978 insert_sample(conn, "s2", "snare.wav", "wav");
979
980 // Tags
981 conn.execute(
982 "INSERT INTO tags (sample_hash, tag) VALUES ('s1', 'drums')",
983 [],
984 )
985 .unwrap();
986 conn.execute(
987 "INSERT INTO tags (sample_hash, tag) VALUES ('s2', 'perc')",
988 [],
989 )
990 .unwrap();
991
992 // Collection
993 let now = chrono::Utc::now().timestamp();
994 conn.execute(
995 "INSERT INTO collections (name, description, created_at) VALUES ('Kit', NULL, ?1)",
996 [now],
997 )
998 .unwrap();
999 let coll_id = conn.last_insert_rowid();
1000
1001 // Collection members
1002 conn.execute(
1003 "INSERT INTO collection_members (collection_id, sample_hash, added_at) VALUES (?1, 's1', ?2)",
1004 rusqlite::params![coll_id, now],
1005 ).unwrap();
1006
1007 set_sync_state(conn, "applying_remote", "0").unwrap();
1008 clear_changelog(conn);
1009
1010 let total = create_initial_snapshot(conn).unwrap();
1011 // 2 samples + 2 tags + 1 collection + 1 collection_member
1012 // + 1 user_config seed (sample_tombstone_retain_days, from M019) = 7
1013 assert_eq!(total, 7);
1014
1015 assert_eq!(changelog_count(conn, Some("samples"), Some("INSERT")), 2);
1016 assert_eq!(changelog_count(conn, Some("tags"), Some("INSERT")), 2);
1017 assert_eq!(
1018 changelog_count(conn, Some("collections"), Some("INSERT")),
1019 1
1020 );
1021 assert_eq!(
1022 changelog_count(conn, Some("collection_members"), Some("INSERT")),
1023 1
1024 );
1025 }
1026
1027 #[test]
1028 fn snapshot_changelog_entries_contain_valid_json() {
1029 let db = setup_test_db();
1030 let conn = db.conn();
1031
1032 set_sync_state(conn, "applying_remote", "1").unwrap();
1033 insert_sample(conn, "json_test", "pad.wav", "wav");
1034 set_sync_state(conn, "applying_remote", "0").unwrap();
1035 clear_changelog(conn);
1036
1037 create_initial_snapshot(conn).unwrap();
1038
1039 let data: String = conn
1040 .query_row(
1041 "SELECT data FROM sync_changelog WHERE table_name = 'samples' AND row_id = 'json_test'",
1042 [],
1043 |row| row.get(0),
1044 )
1045 .unwrap();
1046
1047 let parsed: serde_json::Value = serde_json::from_str(&data).unwrap();
1048 assert_eq!(parsed["hash"], "json_test");
1049 assert_eq!(parsed["original_name"], "pad.wav");
1050 assert_eq!(parsed["file_extension"], "wav");
1051 assert!(parsed.get("file_size").is_some());
1052 assert!(parsed.get("import_date").is_some());
1053 // Columns added in migrations 008-009 must be present in snapshot
1054 assert!(
1055 parsed.get("cloud_only").is_some(),
1056 "snapshot missing cloud_only"
1057 );
1058 assert!(
1059 parsed.get("duration").is_some(),
1060 "snapshot missing duration"
1061 );
1062 }
1063
1064 #[test]
1065 fn snapshot_samples_columns_match_whitelist() {
1066 let db = setup_test_db();
1067 let conn = db.conn();
1068
1069 set_sync_state(conn, "applying_remote", "1").unwrap();
1070 insert_sample(conn, "col_test", "check.wav", "wav");
1071 set_sync_state(conn, "applying_remote", "0").unwrap();
1072 clear_changelog(conn);
1073
1074 create_initial_snapshot(conn).unwrap();
1075
1076 let data: String = conn
1077 .query_row(
1078 "SELECT data FROM sync_changelog WHERE table_name = 'samples' AND row_id = 'col_test'",
1079 [],
1080 |row| row.get(0),
1081 )
1082 .unwrap();
1083 let parsed: serde_json::Value = serde_json::from_str(&data).unwrap();
1084 let obj = parsed.as_object().unwrap();
1085
1086 // Every column in the whitelist must appear in the snapshot JSON
1087 for col in table_columns("samples").unwrap() {
1088 assert!(
1089 obj.contains_key(*col),
1090 "snapshot samples missing column: {col}"
1091 );
1092 }
1093 }
1094
1095 #[test]
1096 fn snapshot_audio_analysis_columns_match_whitelist() {
1097 let db = setup_test_db();
1098 let conn = db.conn();
1099
1100 set_sync_state(conn, "applying_remote", "1").unwrap();
1101 insert_sample(conn, "aa_test", "tone.wav", "wav");
1102 conn.execute(
1103 "INSERT INTO audio_analysis (hash, duration, sample_rate, channels, analyzed_at) VALUES ('aa_test', 1.5, 44100, 2, 1_000_000)",
1104 [],
1105 ).unwrap();
1106 set_sync_state(conn, "applying_remote", "0").unwrap();
1107 clear_changelog(conn);
1108
1109 create_initial_snapshot(conn).unwrap();
1110
1111 let data: String = conn.query_row(
1112 "SELECT data FROM sync_changelog WHERE table_name = 'audio_analysis' AND row_id = 'aa_test'",
1113 [],
1114 |row| row.get(0),
1115 ).unwrap();
1116 let parsed: serde_json::Value = serde_json::from_str(&data).unwrap();
1117 let obj = parsed.as_object().unwrap();
1118
1119 for col in table_columns("audio_analysis").unwrap() {
1120 assert!(
1121 obj.contains_key(*col),
1122 "snapshot audio_analysis missing column: {col}"
1123 );
1124 }
1125 }
1126
1127 #[test]
1128 fn loose_files_excluded_from_sync() {
1129 let db = setup_test_db();
1130 let conn = db.conn();
1131 clear_changelog(conn);
1132
1133 // Insert loose_files, should NOT fire trigger
1134 conn.execute(
1135 "INSERT INTO user_config (key, value) VALUES ('loose_files', '1')",
1136 [],
1137 )
1138 .unwrap();
1139 assert_eq!(changelog_count(conn, Some("user_config"), None), 0);
1140
1141 // Update loose_files, should NOT fire trigger
1142 conn.execute(
1143 "UPDATE user_config SET value = '0' WHERE key = 'loose_files'",
1144 [],
1145 )
1146 .unwrap();
1147 assert_eq!(changelog_count(conn, Some("user_config"), None), 0);
1148
1149 // Normal key, should fire trigger
1150 conn.execute(
1151 "INSERT INTO user_config (key, value) VALUES ('theme', 'dark')",
1152 [],
1153 )
1154 .unwrap();
1155 assert_eq!(changelog_count(conn, Some("user_config"), None), 1);
1156 }
1157
1158 #[test]
1159 fn snapshot_after_adding_data_does_not_duplicate() {
1160 let db = setup_test_db();
1161 let conn = db.conn();
1162
1163 set_sync_state(conn, "applying_remote", "1").unwrap();
1164 insert_sample(conn, "first", "one.wav", "wav");
1165 set_sync_state(conn, "applying_remote", "0").unwrap();
1166 clear_changelog(conn);
1167
1168 let first = create_initial_snapshot(conn).unwrap();
1169 // 1 sample + 1 user_config seed (sample_tombstone_retain_days, M019)
1170 assert_eq!(first, 2);
1171
1172 // Add more data after snapshot (these go through normal triggers)
1173 insert_sample(conn, "second", "two.wav", "wav");
1174
1175 // Second snapshot call should return 0 (flag already set)
1176 let second = create_initial_snapshot(conn).unwrap();
1177 assert_eq!(second, 0);
1178
1179 // Total changelog: 1 from snapshot + 1 from trigger
1180 assert_eq!(changelog_count(conn, Some("samples"), None), 2);
1181 }
1182
1183 #[test]
1184 fn cleanup_preserves_recent_pushed_entries() {
1185 let db = setup_test_db();
1186 let conn = db.conn();
1187 clear_changelog(conn);
1188
1189 // Entry pushed exactly 6 days ago (within 7-day window)
1190 conn.execute(
1191 "INSERT INTO sync_changelog (table_name, op, row_id, timestamp, pushed) \
1192 VALUES ('samples', 'INSERT', 'six_days', datetime('now', '-6 days'), 1)",
1193 [],
1194 )
1195 .unwrap();
1196
1197 // Entry pushed exactly 8 days ago (outside 7-day window)
1198 conn.execute(
1199 "INSERT INTO sync_changelog (table_name, op, row_id, timestamp, pushed) \
1200 VALUES ('samples', 'INSERT', 'eight_days', datetime('now', '-8 days'), 1)",
1201 [],
1202 )
1203 .unwrap();
1204
1205 let deleted = cleanup_changelog(conn).unwrap();
1206 assert_eq!(deleted, 1);
1207
1208 // six_days should remain, eight_days should be deleted
1209 let remaining_id: String = conn
1210 .query_row("SELECT row_id FROM sync_changelog", [], |r| r.get(0))
1211 .unwrap();
1212 assert_eq!(remaining_id, "six_days");
1213 }
1214
1215 #[test]
1216 fn enforce_retention_at_exact_cap_is_noop() {
1217 let db = setup_test_db();
1218 let conn = db.conn();
1219 clear_changelog(conn);
1220
1221 for i in 0..MAX_CHANGELOG_ENTRIES {
1222 conn.execute(
1223 "INSERT INTO sync_changelog (table_name, op, row_id, data) \
1224 VALUES ('samples', 'UPDATE', ?1, '{}')",
1225 [format!("row-{i}")],
1226 )
1227 .unwrap();
1228 }
1229
1230 let outcome = enforce_changelog_retention(conn).unwrap();
1231 assert_eq!(outcome.deleted, 0);
1232
1233 let count: i64 = conn
1234 .query_row("SELECT COUNT(*) FROM sync_changelog", [], |r| r.get(0))
1235 .unwrap();
1236 assert_eq!(count, MAX_CHANGELOG_ENTRIES);
1237 }
1238
1239 #[test]
1240 fn cleanup_and_retention_combined() {
1241 let db = setup_test_db();
1242 let conn = db.conn();
1243 clear_changelog(conn);
1244
1245 // Insert many old pushed entries (should be cleaned by cleanup)
1246 for i in 0..100 {
1247 conn.execute(
1248 "INSERT INTO sync_changelog (table_name, op, row_id, timestamp, pushed) \
1249 VALUES ('samples', 'UPDATE', ?1, datetime('now', '-14 days'), 1)",
1250 [format!("old-{i}")],
1251 )
1252 .unwrap();
1253 }
1254
1255 // Insert entries at the retention cap to test combined behavior
1256 for i in 0..MAX_CHANGELOG_ENTRIES {
1257 conn.execute(
1258 "INSERT INTO sync_changelog (table_name, op, row_id, data) \
1259 VALUES ('samples', 'INSERT', ?1, '{}')",
1260 [format!("new-{i}")],
1261 )
1262 .unwrap();
1263 }
1264
1265 // cleanup removes old pushed entries
1266 let cleaned = cleanup_changelog(conn).unwrap();
1267 assert_eq!(cleaned, 100);
1268
1269 // retention should be a noop now (exactly MAX_CHANGELOG_ENTRIES remain)
1270 let retained = enforce_changelog_retention(conn).unwrap();
1271 assert_eq!(retained.deleted, 0);
1272 }
1273
1274 #[test]
1275 fn sync_state_get_and_set_roundtrip() {
1276 let db = setup_test_db();
1277 let conn = db.conn();
1278
1279 set_sync_state(conn, "auto_sync_enabled", "1").unwrap();
1280 assert_eq!(get_sync_state(conn, "auto_sync_enabled").unwrap(), "1");
1281
1282 set_sync_state(conn, "auto_sync_enabled", "0").unwrap();
1283 assert_eq!(get_sync_state(conn, "auto_sync_enabled").unwrap(), "0");
1284 }
1285
1286 // Download query logic
1287
1288 #[test]
1289 fn missing_blobs_query_finds_sync_enabled_only() {
1290 let db = setup_test_db();
1291 let conn = db.conn();
1292
1293 // Create two VFS: one with sync_files=true, one without
1294 let vfs_sync = insert_vfs(conn, "Synced", true);
1295 let vfs_local = insert_vfs(conn, "Local", false);
1296
1297 insert_sample(conn, "hash_a", "kick.wav", "wav");
1298 insert_sample(conn, "hash_b", "snare.wav", "wav");
1299
1300 // Link hash_a to synced VFS, hash_b to local-only VFS
1301 let now = chrono::Utc::now().timestamp();
1302 conn.execute(
1303 "INSERT INTO vfs_nodes (vfs_id, parent_id, name, node_type, sample_hash, created_at) VALUES (?1, NULL, 'kick.wav', 'sample', 'hash_a', ?2)",
1304 rusqlite::params![vfs_sync, now],
1305 ).unwrap();
1306 conn.execute(
1307 "INSERT INTO vfs_nodes (vfs_id, parent_id, name, node_type, sample_hash, created_at) VALUES (?1, NULL, 'snare.wav', 'sample', 'hash_b', ?2)",
1308 rusqlite::params![vfs_local, now],
1309 ).unwrap();
1310
1311 // Run the same query as download_missing_blobs
1312 let mut stmt = conn
1313 .prepare(
1314 "SELECT DISTINCT s.hash, s.file_extension
1315 FROM samples s
1316 JOIN vfs_nodes vn ON vn.sample_hash = s.hash
1317 JOIN vfs v ON v.id = vn.vfs_id
1318 WHERE v.sync_files = 1",
1319 )
1320 .unwrap();
1321 let rows: Vec<(String, String)> = stmt
1322 .query_map([], |row| Ok((row.get(0)?, row.get(1)?)))
1323 .unwrap()
1324 .collect::<std::result::Result<Vec<_>, _>>()
1325 .unwrap();
1326
1327 assert_eq!(rows.len(), 1);
1328 assert_eq!(rows[0].0, "hash_a");
1329 }
1330
1331 #[test]
1332 fn missing_blobs_query_deduplicates_multi_vfs() {
1333 let db = setup_test_db();
1334 let conn = db.conn();
1335
1336 let vfs1 = insert_vfs(conn, "Lib1", true);
1337 let vfs2 = insert_vfs(conn, "Lib2", true);
1338 insert_sample(conn, "shared_hash", "shared.wav", "wav");
1339
1340 // Same sample linked in two synced VFS entries
1341 let now = chrono::Utc::now().timestamp();
1342 conn.execute(
1343 "INSERT INTO vfs_nodes (vfs_id, parent_id, name, node_type, sample_hash, created_at) VALUES (?1, NULL, 'a.wav', 'sample', 'shared_hash', ?2)",
1344 rusqlite::params![vfs1, now],
1345 ).unwrap();
1346 conn.execute(
1347 "INSERT INTO vfs_nodes (vfs_id, parent_id, name, node_type, sample_hash, created_at) VALUES (?1, NULL, 'b.wav', 'sample', 'shared_hash', ?2)",
1348 rusqlite::params![vfs2, now],
1349 ).unwrap();
1350
1351 let mut stmt = conn
1352 .prepare(
1353 "SELECT DISTINCT s.hash, s.file_extension
1354 FROM samples s
1355 JOIN vfs_nodes vn ON vn.sample_hash = s.hash
1356 JOIN vfs v ON v.id = vn.vfs_id
1357 WHERE v.sync_files = 1",
1358 )
1359 .unwrap();
1360 let rows: Vec<(String, String)> = stmt
1361 .query_map([], |row| Ok((row.get(0)?, row.get(1)?)))
1362 .unwrap()
1363 .collect::<std::result::Result<Vec<_>, _>>()
1364 .unwrap();
1365
1366 // DISTINCT should yield exactly one row
1367 assert_eq!(rows.len(), 1);
1368 }
1369
1370 // Upload query logic
1371
1372 #[test]
1373 fn upload_pending_query_finds_non_cloud_only() {
1374 let db = setup_test_db();
1375 let conn = db.conn();
1376
1377 let vfs = insert_vfs(conn, "Synced", true);
1378 insert_sample(conn, "local_hash", "kick.wav", "wav");
1379
1380 // Mark one sample as cloud_only
1381 conn.execute(
1382 "UPDATE samples SET cloud_only = 1 WHERE hash = 'local_hash'",
1383 [],
1384 )
1385 .unwrap();
1386
1387 insert_sample(conn, "present_hash", "snare.wav", "wav");
1388
1389 let now = chrono::Utc::now().timestamp();
1390 conn.execute(
1391 "INSERT INTO vfs_nodes (vfs_id, parent_id, name, node_type, sample_hash, created_at) VALUES (?1, NULL, 'kick.wav', 'sample', 'local_hash', ?2)",
1392 rusqlite::params![vfs, now],
1393 ).unwrap();
1394 conn.execute(
1395 "INSERT INTO vfs_nodes (vfs_id, parent_id, name, node_type, sample_hash, created_at) VALUES (?1, NULL, 'snare.wav', 'sample', 'present_hash', ?2)",
1396 rusqlite::params![vfs, now],
1397 ).unwrap();
1398
1399 // Same query as upload_pending_blobs
1400 let mut stmt = conn
1401 .prepare(
1402 "SELECT DISTINCT s.hash, s.file_extension, s.file_size
1403 FROM samples s
1404 JOIN vfs_nodes vn ON vn.sample_hash = s.hash
1405 JOIN vfs v ON v.id = vn.vfs_id
1406 WHERE v.sync_files = 1 AND s.cloud_only = 0",
1407 )
1408 .unwrap();
1409 let rows: Vec<(String, String, i64)> = stmt
1410 .query_map([], |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)))
1411 .unwrap()
1412 .collect::<std::result::Result<Vec<_>, _>>()
1413 .unwrap();
1414
1415 // Only present_hash (cloud_only=0)
1416 assert_eq!(rows.len(), 1);
1417 assert_eq!(rows[0].0, "present_hash");
1418 }
1419
1420 #[test]
1421 fn push_changelog_reads_unpushed_in_order() {
1422 let db = setup_test_db();
1423 let conn = db.conn();
1424 clear_changelog(conn);
1425
1426 // Insert changelog entries with specific order
1427 for i in 0..5 {
1428 conn.execute(
1429 "INSERT INTO sync_changelog (table_name, op, row_id, data, pushed) VALUES ('samples', 'INSERT', ?1, '{}', 0)",
1430 [format!("push-{i}")],
1431 ).unwrap();
1432 }
1433 // Mark one as already pushed
1434 conn.execute(
1435 "UPDATE sync_changelog SET pushed = 1 WHERE row_id = 'push-2'",
1436 [],
1437 )
1438 .unwrap();
1439
1440 // Same query as push_changes
1441 let mut stmt = conn
1442 .prepare(
1443 "SELECT id, table_name, op, row_id, timestamp, data
1444 FROM sync_changelog
1445 WHERE pushed = 0
1446 ORDER BY id ASC
1447 LIMIT ?1",
1448 )
1449 .unwrap();
1450 let rows: Vec<(i64, String, String, String)> = stmt
1451 .query_map([500i64], |row| {
1452 Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?))
1453 })
1454 .unwrap()
1455 .collect::<std::result::Result<Vec<_>, _>>()
1456 .unwrap();
1457
1458 // Should skip push-2 (already pushed)
1459 assert_eq!(rows.len(), 4);
1460 let row_ids: Vec<&str> = rows.iter().map(|r| r.3.as_str()).collect();
1461 assert!(!row_ids.contains(&"push-2"));
1462 // Should be in ASC order
1463 assert_eq!(row_ids[0], "push-0");
1464 assert_eq!(row_ids[3], "push-4");
1465 }
1466
1467 // Resolve: mixed upsert/delete ordering
1468
1469 #[test]
1470 fn apply_remote_changes_mixed_ops_correct_order() {
1471 let db = setup_test_db();
1472 let conn = db.conn();
1473 clear_changelog(conn);
1474
1475 // Insert a sample first
1476 let changes_insert = vec![
1477 change(
1478 "samples",
1479 ChangeOp::Insert,
1480 "mix_hash",
1481 Some(json!({
1482 "hash": "mix_hash", "original_name": "mixed.wav",
1483 "file_extension": "wav", "file_size": 1024,
1484 "import_date": 1_000_000, "last_modified": 1_000_000,
1485 "cloud_only": 0
1486 })),
1487 ),
1488 change(
1489 "tags",
1490 ChangeOp::Insert,
1491 "mix_hash:bass",
1492 Some(json!({
1493 "sample_hash": "mix_hash", "tag": "bass"
1494 })),
1495 ),
1496 ];
1497 apply_remote_changes(conn, &changes_insert).unwrap();
1498
1499 // Now delete tag then sample (mixed batch)
1500 let changes_delete = vec![
1501 change("tags", ChangeOp::Delete, "mix_hash:bass", None),
1502 change("samples", ChangeOp::Delete, "mix_hash", None),
1503 ];
1504 let applied = apply_remote_changes(conn, &changes_delete).unwrap();
1505 assert_eq!(applied, 2);
1506
1507 // The sample tombstones (row survives, deleted_at set) rather than
1508 // hard-deleting, so a remote delete can never cascade-wipe local data.
1509 let (sample_count, tombstoned): (i64, i64) = conn
1510 .query_row(
1511 "SELECT COUNT(*), COUNT(deleted_at) FROM samples WHERE hash = 'mix_hash'",
1512 [],
1513 |row| Ok((row.get(0)?, row.get(1)?)),
1514 )
1515 .unwrap();
1516 assert_eq!(sample_count, 1);
1517 assert_eq!(tombstoned, 1);
1518
1519 // The explicit tag delete is a real (non-cascading) hard delete.
1520 let tag_count: i64 = conn
1521 .query_row(
1522 "SELECT COUNT(*) FROM tags WHERE sample_hash = 'mix_hash'",
1523 [],
1524 |row| row.get(0),
1525 )
1526 .unwrap();
1527 assert_eq!(tag_count, 0);
1528 }
1529
1530 #[test]
1531 fn apply_upsert_updates_existing_row() {
1532 let db = setup_test_db();
1533 let conn = db.conn();
1534
1535 // Insert initial sample
1536 apply_upsert(
1537 conn,
1538 "samples",
1539 &json!({
1540 "hash": "upd_hash", "original_name": "old.wav",
1541 "file_extension": "wav", "file_size": 1024,
1542 "import_date": 1_000_000, "last_modified": 1_000_000,
1543 "cloud_only": 0
1544 }),
1545 )
1546 .unwrap();
1547
1548 // Update via upsert (ON CONFLICT DO UPDATE)
1549 apply_upsert(
1550 conn,
1551 "samples",
1552 &json!({
1553 "hash": "upd_hash", "original_name": "new.wav",
1554 "file_extension": "wav", "file_size": 2048,
1555 "import_date": 1_000_000, "last_modified": 2_000_000,
1556 "cloud_only": 0
1557 }),
1558 )
1559 .unwrap();
1560
1561 let name: String = conn
1562 .query_row(
1563 "SELECT original_name FROM samples WHERE hash = 'upd_hash'",
1564 [],
1565 |row| row.get(0),
1566 )
1567 .unwrap();
1568 assert_eq!(name, "new.wav");
1569
1570 let size: i64 = conn
1571 .query_row(
1572 "SELECT file_size FROM samples WHERE hash = 'upd_hash'",
1573 [],
1574 |row| row.get(0),
1575 )
1576 .unwrap();
1577 assert_eq!(size, 2048);
1578
1579 // Should still be only 1 row
1580 let count: i64 = conn
1581 .query_row(
1582 "SELECT COUNT(*) FROM samples WHERE hash = 'upd_hash'",
1583 [],
1584 |row| row.get(0),
1585 )
1586 .unwrap();
1587 assert_eq!(count, 1);
1588 }
1589
1590 #[test]
1591 fn apply_delete_nonexistent_is_noop() {
1592 let db = setup_test_db();
1593 let conn = db.conn();
1594
1595 // Delete a hash that doesn't exist, should succeed (no-op)
1596 let result = apply_delete(conn, "samples", "nonexistent_hash", None);
1597 assert!(result.is_ok());
1598 }
1599
1600 // Phase 5: classifier syncs across a user's own devices
1601
1602 /// Drain everything a device would push: each unsuppressed changelog row as a
1603 /// `ChangeEntry`, in insertion order (apply_remote_changes re-orders FK-safely).
1604 fn drain_changelog(conn: &Connection) -> Vec<ChangeEntry> {
1605 let mut stmt = conn
1606 .prepare("SELECT table_name, op, row_id, data FROM sync_changelog ORDER BY id")
1607 .unwrap();
1608 let rows: Vec<(String, String, String, Option<String>)> = stmt
1609 .query_map([], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?)))
1610 .unwrap()
1611 .collect::<std::result::Result<Vec<_>, _>>()
1612 .unwrap();
1613 rows.into_iter()
1614 .map(|(table, op, row_id, data)| ChangeEntry {
1615 table,
1616 op: match op.as_str() {
1617 "INSERT" => ChangeOp::Insert,
1618 "UPDATE" => ChangeOp::Update,
1619 _ => ChangeOp::Delete,
1620 },
1621 row_id,
1622 timestamp: chrono::Utc::now(),
1623 hlc: synckit_client::Hlc::zero(synckit_client::DeviceId::nil()),
1624 data: data.and_then(|d| serde_json::from_str(&d).ok()),
1625 extra: serde_json::Map::default(),
1626 })
1627 .collect()
1628 }
1629
1630 /// Build a classified library on device A, sync it to a fresh device B, and confirm
1631 /// the whole classifier rebuilds from synced data alone, no audio, no re-analysis.
1632 #[test]
1633 fn classifier_rebuilds_on_fresh_device_from_sync() {
1634 use audiofiles_core::analysis::{exemplar, features::FEATURE_VERSION};
1635 use audiofiles_core::rules::{
1636 self, MatchMode, NewRule, RuleAction, RuleCondition, RuleField, RuleOp,
1637 };
1638
1639 // --- Device A: a sample with a feature vector, a rule, and a policy. ---
1640 let a = setup_test_db();
1641 clear_changelog(a.conn());
1642 insert_sample(a.conn(), "k1", "808 kick.wav", "wav");
1643 let vec_json = serde_json::to_string(&vec![0.1f64; 35]).unwrap();
1644 a.conn()
1645 .execute(
1646 "INSERT INTO sample_features (hash, feat_version, vector, computed_at) VALUES ('k1', ?1, ?2, 0)",
1647 rusqlite::params![FEATURE_VERSION, vec_json],
1648 )
1649 .unwrap();
1650 rules::create_rule(
1651 &a,
1652 NewRule {
1653 name: "kicks".into(),
1654 enabled: true,
1655 priority: None,
1656 match_mode: MatchMode::All,
1657 conditions: vec![RuleCondition {
1658 field: RuleField::Name,
1659 op: RuleOp::Contains,
1660 value: "kick".into(),
1661 }],
1662 actions: vec![RuleAction::AddTag("instrument.drum.kick".into())],
1663 },
1664 )
1665 .unwrap();
1666 exemplar::set_policy(&a, "instrument.drum.kick", 0.4, 0.8).unwrap();
1667 // Apply the rule on A so the tag + 'rule' provenance exist to sync.
1668 rules::apply_all_rules(&a).unwrap();
1669
1670 let changes = drain_changelog(a.conn());
1671
1672 // --- Device B: fresh DB, apply the synced changes (no audio anywhere). ---
1673 let b = setup_test_db();
1674 clear_changelog(b.conn());
1675 apply_remote_changes(b.conn(), &changes).unwrap();
1676
1677 // 1. k-NN index rebuilds from the synced feature vector, no re-analysis.
1678 let index = exemplar::build_index(&b).unwrap();
1679 assert_eq!(
1680 index.len(),
1681 1,
1682 "synced exemplar should rebuild the k-NN index"
1683 );
1684
1685 // 2. Rule synced.
1686 let synced_rules = rules::list_rules(&b).unwrap();
1687 assert_eq!(synced_rules.len(), 1);
1688 assert_eq!(synced_rules[0].name, "kicks");
1689
1690 // 3. Policy synced.
1691 let policy = exemplar::get_policy(&b, "instrument.drum.kick").unwrap();
1692 assert!((policy.auto_threshold - 0.8).abs() < 1e-9);
1693 assert!((policy.review_threshold - 0.4).abs() < 1e-9);
1694
1695 // 4. Tag + provenance synced (rule-sourced, manual-sticky semantics intact).
1696 assert!(
1697 audiofiles_core::tags::get_sample_tags(&b, "k1")
1698 .unwrap()
1699 .contains(&"instrument.drum.kick".to_string())
1700 );
1701 let prov = rules::sample_tag_provenance(&b, "k1").unwrap();
1702 assert!(
1703 prov.iter()
1704 .any(|(t, s, _)| t == "instrument.drum.kick" && s == "rule")
1705 );
1706
1707 // 5. Re-running rules on B (from synced metadata) is consistent, no audio needed.
1708 let changed = rules::apply_all_rules(&b).unwrap();
1709 assert_eq!(
1710 changed, 0,
1711 "synced state already reconciled; re-apply is a no-op"
1712 );
1713 }
1714
1715 /// Import an `.afcl` layer on device A, sync to fresh device B, and confirm the layer,
1716 /// its exemplars (in the k-NN index), and its disabled rule all rebuild, then a layer
1717 /// removal on A propagates the deletes to B. (Phase 7c.)
1718 #[test]
1719 fn imported_classifier_layer_syncs_to_fresh_device() {
1720 use audiofiles_core::analysis::afcl::{self, Afcl, AfclExemplar, AfclManifest};
1721 use audiofiles_core::analysis::{exemplar, features::FEATURE_VERSION};
1722 use audiofiles_core::rules::{MatchMode, Rule, RuleAction, RuleCondition, RuleField, RuleOp};
1723
1724 let make_afcl = || Afcl {
1725 manifest: AfclManifest {
1726 afcl_version: afcl::AFCL_VERSION,
1727 feat_version: FEATURE_VERSION,
1728 name: "shared kicks".into(),
1729 description: String::new(),
1730 kind: "imported".into(),
1731 created_at: 0,
1732 license_note: String::new(),
1733 exemplar_count: 2,
1734 rule_count: 1,
1735 policy_count: 0,
1736 },
1737 exemplars: vec![
1738 AfclExemplar {
1739 vector: vec![0.0; 35],
1740 tags: vec!["instrument.drum.kick".into()],
1741 },
1742 AfclExemplar {
1743 vector: vec![0.05; 35],
1744 tags: vec!["instrument.drum.kick".into()],
1745 },
1746 ],
1747 rules: vec![Rule {
1748 id: String::new(),
1749 name: "shared kick rule".into(),
1750 enabled: true, // import forces disabled
1751 priority: 0,
1752 match_mode: MatchMode::All,
1753 conditions: vec![RuleCondition {
1754 field: RuleField::Name,
1755 op: RuleOp::Contains,
1756 value: "kick".into(),
1757 }],
1758 actions: vec![RuleAction::AddTag("instrument.drum.kick".into())],
1759 created_at: 0,
1760 }],
1761 policy: vec![],
1762 };
1763
1764 // --- Device A: import the layer, then drain its changelog. ---
1765 let a = setup_test_db();
1766 clear_changelog(a.conn());
1767 let summary = afcl::import(&a, &make_afcl(), Some("shared.afcl")).unwrap();
1768 let changes = drain_changelog(a.conn());
1769
1770 // --- Device B: fresh DB, apply the synced changes. ---
1771 let b = setup_test_db();
1772 clear_changelog(b.conn());
1773 apply_remote_changes(b.conn(), &changes).unwrap();
1774
1775 // Layer, exemplars, and membership rebuilt.
1776 let layers = afcl::list_layers(&b).unwrap();
1777 assert_eq!(layers.len(), 1);
1778 assert_eq!(layers[0].id, summary.layer_id);
1779 assert_eq!(layers[0].exemplar_count, 2);
1780 assert_eq!(layers[0].rule_count, 1);
1781 // The imported exemplars join B's k-NN index.
1782 assert_eq!(afcl::enabled_imported_exemplars(&b).unwrap().len(), 2);
1783 assert_eq!(exemplar::build_index(&b).unwrap().len(), 2);
1784 // The imported rule synced and is disabled.
1785 let synced = audiofiles_core::rules::list_rules(&b).unwrap();
1786 assert_eq!(synced.len(), 1);
1787 assert!(
1788 !synced[0].enabled,
1789 "imported rule stays disabled across sync"
1790 );
1791
1792 // --- Removal on A propagates to B (explicit child deletes are sync-logged). ---
1793 clear_changelog(a.conn());
1794 afcl::remove_layer(&a, &summary.layer_id).unwrap();
1795 let del_changes = drain_changelog(a.conn());
1796 apply_remote_changes(b.conn(), &del_changes).unwrap();
1797
1798 assert!(
1799 afcl::list_layers(&b).unwrap().is_empty(),
1800 "layer removed on B"
1801 );
1802 assert!(afcl::enabled_imported_exemplars(&b).unwrap().is_empty());
1803 assert!(
1804 audiofiles_core::rules::list_rules(&b).unwrap().is_empty(),
1805 "the layer's imported rule was removed on B too"
1806 );
1807 }
1808