From 465b8068bedf0a1a648061cf35be557aca9954ff Mon Sep 17 00:00:00 2001 From: Emil Rossing Date: Thu, 1 Oct 2026 18:05:29 +0200 Subject: [PATCH] perf(desktop): index deleted observations --- desktop/src-tauri/src/data_export.rs | 107 ++++--- desktop/src-tauri/src/lib.rs | 293 ++++++++++++------ desktop/src-tauri/src/local_api/export.rs | 20 +- desktop/src-tauri/src/local_api/tests.rs | 20 +- desktop/src-tauri/src/observation_index.rs | 18 +- desktop/src-tauri/src/observation_query.rs | 24 +- desktop/src-tauri/src/sync_engine/mod.rs | 2 +- .../lib/__tests__/formPreviewBridge.test.ts | 33 ++ desktop/src/lib/formPreviewBridge.ts | 9 +- desktop/src/lib/tauriClient.ts | 2 + desktop/src/pages/ObservationsPage.tsx | 8 +- .../src/services/synk/GeneratedSyncGateway.ts | 2 +- desktop/src/types/domain.ts | 2 + .../observation-query/src/compiler.test.ts | 48 +++ packages/observation-query/src/compiler.ts | 13 +- 15 files changed, 414 insertions(+), 187 deletions(-) diff --git a/desktop/src-tauri/src/data_export.rs b/desktop/src-tauri/src/data_export.rs index 89ba400d8..55f4081b6 100644 --- a/desktop/src-tauri/src/data_export.rs +++ b/desktop/src-tauri/src/data_export.rs @@ -238,10 +238,6 @@ fn infer_col_kind(values: &[&Value]) -> ColKind { } } -fn extras_deleted(extras: &Option) -> bool { - extras.as_ref().and_then(|e| e.deleted).unwrap_or(false) -} - struct DbExportCandidate { id: String, payload: Value, @@ -253,13 +249,7 @@ struct DbExportCandidate { extras: Option, } -fn row_from_db(c: DbExportCandidate) -> Option { - if c.sync_status == "conflict" { - return None; - } - if extras_deleted(&c.extras) { - return None; - } +fn row_from_db(c: DbExportCandidate) -> ExportRow { let form_type = c .form_type .filter(|s| !s.trim().is_empty()) @@ -297,7 +287,7 @@ fn row_from_db(c: DbExportCandidate) -> Option { .into_iter() .map(|(k, v)| (k, json_value_to_export_cell(&v))) .collect(); - Some(ExportRow { + ExportRow { observation_id: c.id, form_type, form_version, @@ -313,7 +303,7 @@ fn row_from_db(c: DbExportCandidate) -> Option { pending, data, raw_payload: c.payload, - }) + } } /// Load exportable rows; `form_types` empty means all forms. @@ -322,23 +312,25 @@ pub(crate) fn load_export_rows( include_pending: bool, form_types: &[String], ) -> Result, CustodianError> { - let sql = if include_pending { - "SELECT id, payload, form_type, updated_at, dirty, sync_status, last_saved_at, observation_extras - FROM observations - WHERE sync_status != 'conflict' - ORDER BY COALESCE(form_type, ''), id" - } else { + let mut sql = String::from( "SELECT id, payload, form_type, updated_at, dirty, sync_status, last_saved_at, observation_extras FROM observations - WHERE sync_status != 'conflict' AND dirty = 0 - ORDER BY COALESCE(form_type, ''), id" - }; - let mut stmt = conn.prepare(sql)?; - let rows = stmt.query_map([], |row| { + WHERE deleted = 0 AND sync_status != 'conflict'", + ); + if !include_pending { + sql.push_str(" AND dirty = 0"); + } + if !form_types.is_empty() { + let placeholders = vec!["?"; form_types.len()].join(", "); + sql.push_str(&format!(" AND form_type IN ({placeholders})")); + } + sql.push_str(" ORDER BY COALESCE(form_type, ''), id"); + + let mut stmt = conn.prepare(&sql)?; + let rows = stmt.query_map(rusqlite::params_from_iter(form_types.iter()), |row| { let payload_raw: String = row.get(1)?; let payload = serde_json::from_str::(&payload_raw).unwrap_or(Value::Null); let dirty: i64 = row.get(4)?; - let sync_status: String = row.get(5)?; let extras_raw: Option = row.get(7)?; let extras = extras_raw.and_then(|s| { let t = s.trim(); @@ -354,7 +346,7 @@ pub(crate) fn load_export_rows( row.get::<_, Option>(2)?, row.get::<_, Option>(3)?, dirty == 1, - sync_status, + row.get::<_, String>(5)?, row.get::<_, String>(6)?, extras, )) @@ -363,14 +355,7 @@ pub(crate) fn load_export_rows( let mut out = Vec::new(); for row in rows { let (id, payload, form_type, updated_at, dirty, sync_status, last_saved_at, extras) = row?; - if !form_types.is_empty() - && !form_type - .as_deref() - .is_some_and(|ft| form_types.iter().any(|t| t == ft)) - { - continue; - } - if let Some(r) = row_from_db(DbExportCandidate { + out.push(row_from_db(DbExportCandidate { id, payload, form_type, @@ -379,9 +364,7 @@ pub(crate) fn load_export_rows( sync_status, last_saved_at, extras, - }) { - out.push(r); - } + })); } Ok(out) } @@ -1000,7 +983,8 @@ mod tests { conflict_payload TEXT, last_saved_at TEXT NOT NULL, last_pushed_at TEXT, - observation_extras TEXT + observation_extras TEXT, + deleted INTEGER NOT NULL DEFAULT 0 ); "#, ) @@ -1027,6 +1011,53 @@ mod tests { assert_eq!(infer_col_kind(&[&e, &f]), ColKind::Boolean); } + #[test] + fn load_export_rows_filters_in_sql() { + let conn = setup_db(); + for (id, form_type, deleted, dirty, sync_status) in [ + ("requested-clean", "person", 0, 0, "clean"), + ("requested-pending", "person", 0, 1, "dirty"), + ("requested-deleted", "person", 1, 0, "clean"), + ("requested-conflict", "person", 0, 1, "conflict"), + ("unrequested-clean", "household", 0, 0, "clean"), + ] { + conn.execute( + "INSERT INTO observations ( + id, payload, form_type, updated_at, dirty, sync_status, + last_saved_at, observation_extras, deleted + ) VALUES (?1, '{}', ?2, '2026-01-01T00:00:00Z', ?3, ?4, + '2026-01-01T00:00:00Z', '{}', ?5)", + params![id, form_type, dirty, sync_status, deleted], + ) + .unwrap(); + } + + let requested = vec!["person".to_string()]; + let rows = load_export_rows(&conn, false, &requested).unwrap(); + assert_eq!( + rows.iter() + .map(|row| row.observation_id.as_str()) + .collect::>(), + ["requested-clean"] + ); + + let rows = load_export_rows(&conn, true, &requested).unwrap(); + assert_eq!( + rows.iter() + .map(|row| row.observation_id.as_str()) + .collect::>(), + ["requested-clean", "requested-pending"] + ); + + let rows = load_export_rows(&conn, true, &[]).unwrap(); + assert_eq!( + rows.iter() + .map(|row| row.observation_id.as_str()) + .collect::>(), + ["unrequested-clean", "requested-clean", "requested-pending"] + ); + } + #[test] fn flatten_and_write_parquet_roundtrip_shape() { let conn = setup_db(); diff --git a/desktop/src-tauri/src/lib.rs b/desktop/src-tauri/src/lib.rs index 5d83dfc48..9f939f9e5 100644 --- a/desktop/src-tauri/src/lib.rs +++ b/desktop/src-tauri/src/lib.rs @@ -649,6 +649,7 @@ struct ObservationRecord { has_conflict_copy: bool, last_saved_at: String, last_pushed_at: Option, + deleted: bool, extras: Option, } @@ -795,7 +796,8 @@ fn init_db(conn: &Connection) -> Result<(), CustodianError> { conflict_payload TEXT, last_saved_at TEXT NOT NULL, last_pushed_at TEXT, - observation_extras TEXT + observation_extras TEXT, + deleted INTEGER NOT NULL DEFAULT 0 ); CREATE TABLE IF NOT EXISTS observation_history ( backup_id TEXT PRIMARY KEY, @@ -815,6 +817,7 @@ fn init_db(conn: &Connection) -> Result<(), CustodianError> { INSERT OR IGNORE INTO sync_state(id, last_pull_at, last_push_at, last_error) VALUES (1, NULL, NULL, NULL); "#, )?; + migrate_observation_columns(conn)?; migrate_sync_state_columns(conn)?; conn.execute( "INSERT OR IGNORE INTO sync_state(id, last_pull_at, last_push_at, last_error, repository_generation, observation_sync_version, last_attachment_version) VALUES (1, NULL, NULL, NULL, 0, 0, 0)", @@ -876,6 +879,37 @@ fn migrate_repository_generation_fresh_install_defaults( Ok(()) } +fn migrate_observation_columns(conn: &Connection) -> Result<(), rusqlite::Error> { + let mut stmt = conn.prepare("PRAGMA table_info(observations)")?; + let cols: Vec = stmt + .query_map([], |row| row.get::<_, String>(1))? + .collect::>()?; + drop(stmt); + + if !cols.iter().any(|c| c == "deleted") { + conn.execute( + "ALTER TABLE observations ADD COLUMN deleted INTEGER NOT NULL DEFAULT 0", + [], + )?; + conn.execute( + "UPDATE observations + SET deleted = CASE + WHEN json_valid(observation_extras) = 1 + THEN json_extract(observation_extras, '$.deleted') IS 1 + ELSE 0 + END", + [], + )?; + } + + conn.execute( + "CREATE INDEX IF NOT EXISTS idx_observations_deleted_form_type + ON observations(deleted, form_type)", + [], + )?; + Ok(()) +} + fn migrate_sync_state_columns(conn: &Connection) -> Result<(), rusqlite::Error> { let mut stmt = conn.prepare("PRAGMA table_info(sync_state)")?; let cols: Vec = stmt @@ -2239,6 +2273,10 @@ fn upsert_observation_from_api( let payload = serde_json::to_string(&incoming.data)?; let timestamp = now_iso(); let extras_json = serialize_observation_extras(&incoming.extras)?; + let incoming_deleted = incoming + .extras + .as_ref() + .map(|extras| i64::from(extras.deleted.unwrap_or(false))); if let Some((local_dirty, local_remote_updated_at, local_payload)) = existing { if should_mark_conflict(local_dirty, &local_remote_updated_at, &incoming.updated_at) { @@ -2287,14 +2325,16 @@ fn upsert_observation_from_api( sync_status = 'clean', conflict_payload = NULL, last_saved_at = ?4, - observation_extras = COALESCE(?5, observation_extras) - WHERE id = ?6", + observation_extras = COALESCE(?5, observation_extras), + deleted = COALESCE(?6, deleted) + WHERE id = ?7", params![ payload, incoming.form_type, incoming.updated_at, timestamp, extras_json, + incoming_deleted, incoming.observation_id ], )?; @@ -2305,15 +2345,16 @@ fn upsert_observation_from_api( conn.execute( "INSERT INTO observations ( id, payload, form_type, updated_at, remote_updated_at, - dirty, sync_status, conflict_payload, last_saved_at, last_pushed_at, observation_extras - ) VALUES (?1, ?2, ?3, ?4, ?4, 0, 'clean', NULL, ?5, NULL, ?6)", + dirty, sync_status, conflict_payload, last_saved_at, last_pushed_at, observation_extras, deleted + ) VALUES (?1, ?2, ?3, ?4, ?4, 0, 'clean', NULL, ?5, NULL, ?6, ?7)", params![ incoming.observation_id, payload, incoming.form_type, incoming.updated_at, timestamp, - extras_json + extras_json, + incoming_deleted.unwrap_or(0) ], )?; Ok(false) @@ -2340,6 +2381,10 @@ fn upsert_observation_from_local_import( .filter(|s| !s.trim().is_empty()) .unwrap_or_else(|| timestamp.clone()); let extras_json = serialize_observation_extras(&incoming.extras)?; + let incoming_deleted = incoming + .extras + .as_ref() + .map(|extras| i64::from(extras.deleted.unwrap_or(false))); if existing.is_some() { conn.execute( @@ -2351,14 +2396,16 @@ fn upsert_observation_from_local_import( sync_status = 'dirty', conflict_payload = NULL, last_saved_at = ?4, - observation_extras = COALESCE(?5, observation_extras) - WHERE id = ?6", + observation_extras = COALESCE(?5, observation_extras), + deleted = COALESCE(?6, deleted) + WHERE id = ?7", params![ payload, incoming.form_type, updated, timestamp, extras_json, + incoming_deleted, incoming.observation_id ], )?; @@ -2366,15 +2413,16 @@ fn upsert_observation_from_local_import( conn.execute( "INSERT INTO observations ( id, payload, form_type, updated_at, remote_updated_at, - dirty, sync_status, conflict_payload, last_saved_at, last_pushed_at, observation_extras - ) VALUES (?1, ?2, ?3, ?4, NULL, 1, 'dirty', NULL, ?5, NULL, ?6)", + dirty, sync_status, conflict_payload, last_saved_at, last_pushed_at, observation_extras, deleted + ) VALUES (?1, ?2, ?3, ?4, NULL, 1, 'dirty', NULL, ?5, NULL, ?6, ?7)", params![ incoming.observation_id, payload, incoming.form_type, updated, timestamp, - extras_json + extras_json, + incoming_deleted.unwrap_or(0) ], )?; } @@ -2716,11 +2764,16 @@ fn save_observation( .map(serde_json::to_string) .transpose() .map_err(|err| err.to_string())?; + let deleted = req + .extras + .as_ref() + .and_then(|extras| extras.deleted) + .unwrap_or(false); tx.execute( "INSERT INTO observations ( - id, payload, form_type, updated_at, remote_updated_at, dirty, sync_status, conflict_payload, last_saved_at, last_pushed_at, observation_extras - ) VALUES (?1, ?2, ?3, ?4, NULL, 1, 'dirty', NULL, ?5, NULL, ?6) + id, payload, form_type, updated_at, remote_updated_at, dirty, sync_status, conflict_payload, last_saved_at, last_pushed_at, observation_extras, deleted + ) VALUES (?1, ?2, ?3, ?4, NULL, 1, 'dirty', NULL, ?5, NULL, ?6, ?7) ON CONFLICT(id) DO UPDATE SET payload = excluded.payload, form_type = COALESCE(excluded.form_type, observations.form_type), @@ -2729,8 +2782,17 @@ fn save_observation( sync_status = 'dirty', conflict_payload = NULL, last_saved_at = excluded.last_saved_at, - observation_extras = excluded.observation_extras", - params![req.id, payload_raw, req.form_type, logical_updated, timestamp, extras_json], + observation_extras = excluded.observation_extras, + deleted = excluded.deleted", + params![ + req.id, + payload_raw, + req.form_type, + logical_updated, + timestamp, + extras_json, + i64::from(deleted) + ], ) .map_err(|err| err.to_string())?; @@ -2738,8 +2800,7 @@ fn save_observation( let defs = load_active_index_defs(&ctx); if !defs.is_empty() { let ft = req.form_type.as_deref().unwrap_or(""); - let is_deleted = req.extras.as_ref().and_then(|e| e.deleted).unwrap_or(false); - if is_deleted { + if deleted { let _ = observation_index::delete_observation_indexes(&conn, &req.id); } else { let _ = @@ -2758,7 +2819,7 @@ fn get_observation( let conn = open_db(&ctx).map_err(|err| err.to_string())?; let record = conn .query_row( - "SELECT id, payload, form_type, updated_at, remote_updated_at, dirty, sync_status, conflict_payload, last_saved_at, last_pushed_at, observation_extras + "SELECT id, payload, form_type, updated_at, remote_updated_at, dirty, sync_status, conflict_payload, last_saved_at, last_pushed_at, observation_extras, deleted FROM observations WHERE id = ?1", params![id], |row| { @@ -2778,6 +2839,7 @@ fn get_observation( has_conflict_copy: conflict_payload.is_some(), last_saved_at: row.get(8)?, last_pushed_at: row.get(9)?, + deleted: row.get::<_, i64>(11)? == 1, extras: parse_observation_extras(extras_raw), }) }, @@ -2801,7 +2863,7 @@ fn list_observations( let mut stmt = conn .prepare( - "SELECT id, payload, form_type, updated_at, remote_updated_at, dirty, sync_status, conflict_payload, last_saved_at, last_pushed_at, observation_extras + "SELECT id, payload, form_type, updated_at, remote_updated_at, dirty, sync_status, conflict_payload, last_saved_at, last_pushed_at, observation_extras, deleted FROM observations WHERE lower(id) LIKE ?1 OR lower(COALESCE(form_type, '')) LIKE ?1 ORDER BY last_saved_at DESC @@ -2809,26 +2871,7 @@ fn list_observations( ) .map_err(|err| err.to_string())?; let rows = stmt - .query_map(params![pattern, max_rows], |row| { - let payload_raw: String = row.get(1)?; - let payload = serde_json::from_str::(&payload_raw).unwrap_or(Value::Null); - let status: String = row.get(6)?; - let conflict_payload: Option = row.get(7)?; - let extras_raw: Option = row.get(10)?; - Ok(ObservationRecord { - id: row.get(0)?, - payload, - form_type: row.get(2)?, - updated_at: row.get(3)?, - remote_updated_at: row.get(4)?, - dirty: row.get::<_, i64>(5)? == 1, - sync_status: SyncStatus::from(status.as_str()), - has_conflict_copy: conflict_payload.is_some(), - last_saved_at: row.get(8)?, - last_pushed_at: row.get(9)?, - extras: parse_observation_extras(extras_raw), - }) - }) + .query_map(params![pattern, max_rows], map_observation_row) .map_err(|err| err.to_string())?; let mut result = Vec::new(); @@ -2849,6 +2892,7 @@ struct ListObservationsPageResult { fn list_observations_page( query: Option, form_type: Option, + include_deleted: Option, limit: Option, offset: Option, ctx: tauri::State<'_, AppCtxHandle>, @@ -2861,47 +2905,48 @@ fn list_observations_page( let form_filter = form_type .map(|s| s.trim().to_string()) .filter(|s| !s.is_empty()); + let live_filter = if include_deleted.unwrap_or(true) { + "" + } else { + " AND deleted = 0" + }; let total: i64 = if let Some(ref ft) = form_filter { - conn.query_row( + let sql = format!( "SELECT COUNT(*) FROM observations WHERE (lower(id) LIKE ?1 OR lower(COALESCE(form_type, '')) LIKE ?1) - AND COALESCE(form_type, '') = ?2", - params![pattern, ft], - |row| row.get(0), - ) - .map_err(|err| err.to_string())? + AND COALESCE(form_type, '') = ?2{live_filter}" + ); + conn.query_row(&sql, params![pattern, ft], |row| row.get(0)) + .map_err(|err| err.to_string())? } else { - conn.query_row( + let sql = format!( "SELECT COUNT(*) FROM observations - WHERE lower(id) LIKE ?1 OR lower(COALESCE(form_type, '')) LIKE ?1", - params![pattern], - |row| row.get(0), - ) - .map_err(|err| err.to_string())? + WHERE (lower(id) LIKE ?1 OR lower(COALESCE(form_type, '')) LIKE ?1){live_filter}" + ); + conn.query_row(&sql, params![pattern], |row| row.get(0)) + .map_err(|err| err.to_string())? }; - let mut stmt = if form_filter.is_some() { - conn.prepare( - "SELECT id, payload, form_type, updated_at, remote_updated_at, dirty, sync_status, conflict_payload, last_saved_at, last_pushed_at, observation_extras + let sql = if form_filter.is_some() { + format!( + "SELECT id, payload, form_type, updated_at, remote_updated_at, dirty, sync_status, conflict_payload, last_saved_at, last_pushed_at, observation_extras, deleted FROM observations WHERE (lower(id) LIKE ?1 OR lower(COALESCE(form_type, '')) LIKE ?1) - AND COALESCE(form_type, '') = ?2 + AND COALESCE(form_type, '') = ?2{live_filter} ORDER BY last_saved_at DESC - LIMIT ?3 OFFSET ?4", + LIMIT ?3 OFFSET ?4" ) - .map_err(|err| err.to_string())? } else { - conn.prepare( - "SELECT id, payload, form_type, updated_at, remote_updated_at, dirty, sync_status, conflict_payload, last_saved_at, last_pushed_at, observation_extras + format!( + "SELECT id, payload, form_type, updated_at, remote_updated_at, dirty, sync_status, conflict_payload, last_saved_at, last_pushed_at, observation_extras, deleted FROM observations - WHERE lower(id) LIKE ?1 OR lower(COALESCE(form_type, '')) LIKE ?1 + WHERE (lower(id) LIKE ?1 OR lower(COALESCE(form_type, '')) LIKE ?1){live_filter} ORDER BY last_saved_at DESC - LIMIT ?2 OFFSET ?3", + LIMIT ?2 OFFSET ?3" ) - .map_err(|err| err.to_string())? }; - + let mut stmt = conn.prepare(&sql).map_err(|err| err.to_string())?; let rows = if let Some(ref ft) = form_filter { stmt.query_map(params![pattern, ft, max_rows, off], map_observation_row) .map_err(|err| err.to_string())? @@ -2934,6 +2979,7 @@ fn map_observation_row(row: &rusqlite::Row<'_>) -> rusqlite::Result(11)? == 1, extras: parse_observation_extras(extras_raw), }) } @@ -3045,6 +3091,11 @@ fn query_observations( let defs = load_active_index_defs(&ctx); let mut index_keys = observation_index::index_keys_set(&defs); let filter_ref = req.filter.as_ref(); + if req.include_deleted.unwrap_or(false) { + // Tombstones are deliberately absent from observation_index. Including them + // therefore requires payload JSON filtering to preserve query semantics. + index_keys.clear(); + } if filter_ref.is_some() && !index_keys.is_empty() { let active_generation = observation_index::active_generation(&conn).unwrap_or(1); let has_index_rows = conn @@ -3226,7 +3277,7 @@ fn list_dirty_observations( let conn = open_db(&ctx).map_err(|err| err.to_string())?; let mut stmt = conn .prepare( - "SELECT id, payload, form_type, updated_at, remote_updated_at, dirty, sync_status, conflict_payload, last_saved_at, last_pushed_at, observation_extras + "SELECT id, payload, form_type, updated_at, remote_updated_at, dirty, sync_status, conflict_payload, last_saved_at, last_pushed_at, observation_extras, deleted FROM observations WHERE dirty = 1 AND sync_status = 'dirty' ORDER BY last_saved_at ASC @@ -3258,7 +3309,7 @@ pub(crate) fn load_dirty_observations_by_ids( for chunk in ids.chunks(400) { let placeholders = chunk.iter().map(|_| "?").collect::>().join(","); let sql = format!( - "SELECT id, payload, form_type, updated_at, remote_updated_at, dirty, sync_status, conflict_payload, last_saved_at, last_pushed_at, observation_extras + "SELECT id, payload, form_type, updated_at, remote_updated_at, dirty, sync_status, conflict_payload, last_saved_at, last_pushed_at, observation_extras, deleted FROM observations WHERE dirty = 1 AND sync_status = 'dirty' AND id IN ({placeholders})" ); @@ -3488,6 +3539,7 @@ fn build_observation_overview( COUNT(*) AS observation_count, SUM(CASE WHEN dirty = 1 AND sync_status = 'dirty' THEN 1 ELSE 0 END) AS pending_sync_count FROM observations + WHERE deleted = 0 GROUP BY 1 ORDER BY 1 COLLATE NOCASE", )?; @@ -3514,7 +3566,8 @@ fn build_observation_overview( observation_extras, updated_at, last_saved_at - FROM observations", + FROM observations + WHERE deleted = 0", )?; let mut dates: Vec = Vec::new(); @@ -5515,20 +5568,28 @@ fn import_observations_run( } } if !index_defs.is_empty() { - // Sync pull and local file import both update the active generation - // incrementally. A full rebuild is reserved for bundle apply / empty - // index / explicit rebuild — not for adding a few hundred import rows - // on top of an already-indexed sync. - let payload = serde_json::to_string(&observation.data).map_err(|e| e.to_string())?; - let form_type = observation.form_type.as_deref().unwrap_or(""); - observation_index::incremental_reindex( - &tx, - &observation.observation_id, - form_type, - &payload, - &index_defs, - ) - .map_err(|err| err.to_string())?; + // Index the canonical local row: local dirty data may win over an incoming + // pull, and tombstones must never remain in the payload index. + let (payload, form_type, deleted): (String, Option, i64) = tx + .query_row( + "SELECT payload, form_type, deleted FROM observations WHERE id = ?1", + params![observation.observation_id], + |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)), + ) + .map_err(|err| err.to_string())?; + if deleted == 1 { + observation_index::delete_observation_indexes(&tx, &observation.observation_id) + .map_err(|err| err.to_string())?; + } else { + observation_index::incremental_reindex( + &tx, + &observation.observation_id, + form_type.as_deref().unwrap_or(""), + &payload, + &index_defs, + ) + .map_err(|err| err.to_string())?; + } } imported += 1; } @@ -5937,12 +5998,12 @@ mod tests { ObservationExtras, SimpleFileOptions, ZipArchive, ZipWriter, apply_app_bundle_zip_at_workspace, attachment_copy_progress_step, bind_query_params, build_observation_overview, extract_observations_from_json_value, - import_observation_apparently_synced, init_db, mirror_custom_app_dev_folder, - parse_observation_extras, parse_time, publish_bundle_zip_entry_allowed, - resolve_attachment_path, scan_import_json_sync_appearance, - should_emit_attachment_copy_progress, should_mark_conflict, strip_ode_desktop_injection, - upsert_observation_from_local_import, validate_custom_app_dev_source_folder, - zip_dev_mirror_bundle, + import_observation_apparently_synced, init_db, migrate_observation_columns, + mirror_custom_app_dev_folder, parse_observation_extras, parse_time, + publish_bundle_zip_entry_allowed, resolve_attachment_path, + scan_import_json_sync_appearance, should_emit_attachment_copy_progress, + should_mark_conflict, strip_ode_desktop_injection, upsert_observation_from_local_import, + validate_custom_app_dev_source_folder, zip_dev_mirror_bundle, }; use crate::observation_query::SqlParam; use rusqlite::{Connection, params}; @@ -6271,6 +6332,57 @@ mod tests { let _ = fs::remove_dir_all(&base); } + #[test] + fn migrate_observation_columns_backfills_deleted_and_creates_index() { + let conn = Connection::open_in_memory().unwrap(); + conn.execute_batch( + r#" + CREATE TABLE observations ( + id TEXT PRIMARY KEY, + payload TEXT NOT NULL, + form_type TEXT, + observation_extras TEXT + ); + INSERT INTO observations VALUES ('true', '{}', 'person', '{"deleted":true}'); + INSERT INTO observations VALUES ('false', '{}', 'person', '{"deleted":false}'); + INSERT INTO observations VALUES ('missing', '{}', 'person', '{}'); + INSERT INTO observations VALUES ('null', '{}', 'person', NULL); + INSERT INTO observations VALUES ('malformed', '{}', 'person', 'not json'); + "#, + ) + .unwrap(); + + migrate_observation_columns(&conn).unwrap(); + migrate_observation_columns(&conn).unwrap(); + + for (id, expected) in [ + ("true", 1), + ("false", 0), + ("missing", 0), + ("null", 0), + ("malformed", 0), + ] { + let deleted: i64 = conn + .query_row( + "SELECT deleted FROM observations WHERE id = ?1", + params![id], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(deleted, expected, "unexpected deleted value for {id}"); + } + + let index_exists: i64 = conn + .query_row( + "SELECT COUNT(*) FROM sqlite_master + WHERE type = 'index' AND name = 'idx_observations_deleted_form_type'", + [], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(index_exists, 1); + } + #[test] fn upsert_observation_from_local_import_persists_extras() { let conn = Connection::open_in_memory().unwrap(); @@ -6283,6 +6395,7 @@ mod tests { extras: Some(ObservationExtras { author: Some("username:device02".to_string()), tags: Some(vec!["migrated".to_string()]), + deleted: Some(true), geolocation: Some(serde_json::json!({ "latitude": 5.33, "longitude": 36.07 @@ -6291,14 +6404,16 @@ mod tests { }), }; upsert_observation_from_local_import(&conn, &incoming).unwrap(); - let extras_raw: Option = conn + let (extras_raw, deleted): (Option, i64) = conn .query_row( - "SELECT observation_extras FROM observations WHERE id = ?1", + "SELECT observation_extras, deleted FROM observations WHERE id = ?1", params!["uuid:import-1"], - |row| row.get(0), + |row| Ok((row.get(0)?, row.get(1)?)), ) .unwrap(); let parsed = parse_observation_extras(extras_raw).unwrap(); + assert_eq!(deleted, 1); + assert_eq!(parsed.deleted, Some(true)); assert_eq!(parsed.author.as_deref(), Some("username:device02")); assert_eq!(parsed.tags.as_deref(), Some(&["migrated".to_string()][..])); assert!(parsed.geolocation.is_some()); diff --git a/desktop/src-tauri/src/local_api/export.rs b/desktop/src-tauri/src/local_api/export.rs index ee88689c0..0bc8ccfc1 100644 --- a/desktop/src-tauri/src/local_api/export.rs +++ b/desktop/src-tauri/src/local_api/export.rs @@ -4,7 +4,7 @@ use std::collections::{BTreeMap, BTreeSet}; use std::path::{Path, PathBuf}; use std::time::Duration; -use rusqlite::{Connection, OpenFlags}; +use rusqlite::Connection; use serde::Serialize; use super::config::{LocalConfig, workspace_for}; @@ -14,7 +14,7 @@ use super::{ApiError, ApiResult, ErrorCode, hint_lines}; use crate::data_export::{ ExportContext, ExportParquetRequest, ExportProgressFn, load_export_rows, write_parquet_export, }; -use crate::{ServerProfile, sqlite_path_for_workspace}; +use crate::{ServerProfile, init_db, sqlite_path_for_workspace}; /// Manifest field metadata + snippet hint for exports of `profile` (Desktop UI and CLI). pub(crate) fn export_context(profile: &ServerProfile, workspace: &Path) -> ExportContext { @@ -60,7 +60,7 @@ fn io(profile_id: &str, message: impl Into) -> ApiError { ApiError::new(ErrorCode::Io, message).with_profile(profile_id) } -fn open_read_only(workspace: &Path, profile_id: &str) -> ApiResult { +fn open_database(workspace: &Path, profile_id: &str) -> ApiResult { let path = sqlite_path_for_workspace(workspace); if !path.is_file() { return Err(ApiError::new( @@ -69,11 +69,7 @@ fn open_read_only(workspace: &Path, profile_id: &str) -> ApiResult { ) .with_profile(profile_id)); } - let conn = Connection::open_with_flags( - &path, - OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_NO_MUTEX, - ) - .map_err(|e| { + let conn = Connection::open(&path).map_err(|e| { io( profile_id, format!("Could not open the local database: {e}"), @@ -81,6 +77,12 @@ fn open_read_only(workspace: &Path, profile_id: &str) -> ApiResult { })?; conn.busy_timeout(Duration::from_secs(30)) .map_err(|e| io(profile_id, e.to_string()))?; + init_db(&conn).map_err(|e| { + io( + profile_id, + format!("Could not migrate the local database: {e}"), + ) + })?; Ok(conn) } @@ -125,7 +127,7 @@ pub fn export_parquet( progress(0, 1, "Reading observations…"); let rows = { - let conn = open_read_only(&workspace, id)?; + let conn = open_database(&workspace, id)?; load_export_rows(&conn, opts.include_pending, &form_types) .map_err(|e| io(id, e.to_string()))? }; diff --git a/desktop/src-tauri/src/local_api/tests.rs b/desktop/src-tauri/src/local_api/tests.rs index e8ded5b06..0c814629b 100644 --- a/desktop/src-tauri/src/local_api/tests.rs +++ b/desktop/src-tauri/src/local_api/tests.rs @@ -118,7 +118,7 @@ fn disabled_profile_is_rejected_like_unknown() { } /// Seeds the `data` profile: bundle forms `household` + `person`, observations for -/// `household` (2 synced, 1 pending) and `person` (1 synced). +/// `household` (2 synced, 1 pending, 1 deleted) and `person` (1 synced). fn seed_data_profile(base: &Path) { let ws = base.join("data"); write_form(&ws.join("bundles/active/forms"), "household"); @@ -126,21 +126,23 @@ fn seed_data_profile(base: &Path) { fs::create_dir_all(ws.join("sqlite")).unwrap(); let conn = rusqlite::Connection::open(crate::sqlite_path_for_workspace(&ws)).unwrap(); crate::init_db(&conn).unwrap(); - for (id, form, dirty) in [ - ("h1", "household", 0), - ("h2", "household", 0), - ("h3", "household", 1), - ("p1", "person", 0), + for (id, form, dirty, deleted) in [ + ("h1", "household", 0, 0), + ("h2", "household", 0, 0), + ("h3", "household", 1, 0), + ("h-deleted", "household", 0, 1), + ("p1", "person", 0, 0), ] { conn.execute( - "INSERT INTO observations (id, payload, form_type, updated_at, dirty, sync_status, last_saved_at) - VALUES (?1, ?2, ?3, '2026-01-01T00:00:00Z', ?4, ?5, '2026-01-01T00:00:00Z')", + "INSERT INTO observations (id, payload, form_type, updated_at, dirty, sync_status, last_saved_at, deleted) + VALUES (?1, ?2, ?3, '2026-01-01T00:00:00Z', ?4, ?5, '2026-01-01T00:00:00Z', ?6)", rusqlite::params![ id, json!({ "name": id }).to_string(), form, dirty, - if dirty == 1 { "dirty" } else { "clean" } + if dirty == 1 { "dirty" } else { "clean" }, + deleted ], ) .unwrap(); diff --git a/desktop/src-tauri/src/observation_index.rs b/desktop/src-tauri/src/observation_index.rs index af2e74707..5a98eabf8 100644 --- a/desktop/src-tauri/src/observation_index.rs +++ b/desktop/src-tauri/src/observation_index.rs @@ -217,7 +217,11 @@ pub fn rebuild_all_indexes( params![new_gen], )?; - let total: i64 = conn.query_row("SELECT COUNT(*) FROM observations", [], |r| r.get(0))?; + let total: i64 = conn.query_row( + "SELECT COUNT(*) FROM observations WHERE deleted = 0", + [], + |r| r.get(0), + )?; if let Some(ref mut cb) = progress { cb(0, total, Some("Indexing observations…")); } @@ -235,7 +239,8 @@ pub fn rebuild_all_indexes( // 2. Parallel map: parse each payload once and extract all EAV rows. // 3. Join: batch INSERT in one transaction (avoids per-row autocommit). let observations: Vec<(String, String, String)> = { - let mut stmt = conn.prepare("SELECT id, form_type, payload FROM observations")?; + let mut stmt = + conn.prepare("SELECT id, form_type, payload FROM observations WHERE deleted = 0")?; let mapped = stmt.query_map([], |row| { Ok(( row.get::<_, String>(0)?, @@ -433,7 +438,8 @@ mod tests { "CREATE TABLE observations ( id TEXT PRIMARY KEY, form_type TEXT, - payload TEXT NOT NULL + payload TEXT NOT NULL, + deleted INTEGER NOT NULL DEFAULT 0 );", ) .unwrap(); @@ -526,9 +532,9 @@ mod tests { fn rebuild_swaps_generation() { let conn = test_conn(); let defs = sample_defs(); - conn.execute( - "INSERT INTO observations (id, form_type, payload) VALUES ('obs1', 'person', '{\"p_id\":\"P1\"}')", - [], + conn.execute_batch( + "INSERT INTO observations (id, form_type, payload) VALUES ('obs1', 'person', '{\"p_id\":\"P1\"}'); + INSERT INTO observations (id, form_type, payload, deleted) VALUES ('deleted', 'person', '{\"p_id\":\"P2\"}', 1);", ) .unwrap(); let gen1 = rebuild_all_indexes(&conn, &defs, None).unwrap(); diff --git a/desktop/src-tauri/src/observation_query.rs b/desktop/src-tauri/src/observation_query.rs index 3c6e557eb..2915a1555 100644 --- a/desktop/src-tauri/src/observation_query.rs +++ b/desktop/src-tauri/src/observation_query.rs @@ -46,11 +46,7 @@ pub fn compile_observation_query( let mut where_parts = Vec::new(); if !include_deleted { - where_parts.push( - "(json_valid(o.observation_extras) = 0 - OR json_extract(o.observation_extras, '$.deleted') IS NOT 1)" - .to_string(), - ); + where_parts.push("o.deleted = 0".to_string()); } let normalized_form_type = form_type.trim(); @@ -64,9 +60,13 @@ pub fn compile_observation_query( where_parts.push(sql); } + let where_clause = if where_parts.is_empty() { + String::new() + } else { + format!(" WHERE {}", where_parts.join(" AND ")) + }; let sql = format!( - "SELECT o.id, o.payload, o.form_type, o.updated_at, o.remote_updated_at, o.dirty, o.sync_status, o.conflict_payload, o.last_saved_at, o.last_pushed_at, o.observation_extras FROM observations o WHERE {}", - where_parts.join(" AND ") + "SELECT o.id, o.payload, o.form_type, o.updated_at, o.remote_updated_at, o.dirty, o.sync_status, o.conflict_payload, o.last_saved_at, o.last_pushed_at, o.observation_extras, o.deleted FROM observations o{where_clause}" ); Ok(CompiledSql { @@ -402,18 +402,16 @@ mod tests { let live_only = compile_observation_query("person", false, None, &index_keys).unwrap(); - assert!( - live_only - .sql - .contains("json_extract(o.observation_extras, '$.deleted') IS NOT 1") - ); + assert!(live_only.sql.contains("o.deleted = 0")); + assert!(!live_only.sql.contains("json_extract(o.observation_extras")); let include_deleted = compile_observation_query("person", true, None, &index_keys).unwrap(); + assert!(!include_deleted.sql.contains("o.deleted = 0")); assert!( !include_deleted .sql - .contains("json_extract(o.observation_extras, '$.deleted') IS NOT 1") + .contains("json_extract(o.observation_extras") ); } } diff --git a/desktop/src-tauri/src/sync_engine/mod.rs b/desktop/src-tauri/src/sync_engine/mod.rs index a247f0f84..abaec8842 100644 --- a/desktop/src-tauri/src/sync_engine/mod.rs +++ b/desktop/src-tauri/src/sync_engine/mod.rs @@ -225,7 +225,7 @@ fn observation_to_push_json(o: &crate::ObservationRecord) -> Value { "data": o.payload.clone(), "created_at": created_at, "updated_at": updated_at, - "deleted": ex.and_then(|e| e.deleted).unwrap_or(false), + "deleted": o.deleted, "synced_at": ex.and_then(|e| e.synced_at.clone()), "geolocation": ex.and_then(|e| e.geolocation.clone()), "author": ex.and_then(|e| e.author.clone()), diff --git a/desktop/src/lib/__tests__/formPreviewBridge.test.ts b/desktop/src/lib/__tests__/formPreviewBridge.test.ts index 024ccf899..47c63fb38 100644 --- a/desktop/src/lib/__tests__/formPreviewBridge.test.ts +++ b/desktop/src/lib/__tests__/formPreviewBridge.test.ts @@ -122,6 +122,39 @@ describe('handleFormPreviewBridgeMessage', () => { expect(payload.result).toBe(FORM_PREVIEW_FORMULUS_INTERFACE_VERSION); }); + it.each([ + [false, undefined], + [true, true], + ] as const)( + 'passes includeDeleted=%s to native getObservations pagination', + async (expected, requested) => { + const listObservationsPage = vi.mocked(tauriClient.listObservationsPage); + listObservationsPage.mockClear(); + const postMessage = vi.fn(); + const iframe = { + contentWindow: { postMessage } as unknown as Window, + } as HTMLIFrameElement; + + await handleFormPreviewBridgeMessage( + bridgeMessageFromIframe(iframe, { + type: 'getObservations', + messageId: `observations-${expected}`, + formType: 'demo', + ...(requested === undefined ? {} : { includeDeleted: requested }), + }), + { iframe, onFinalize: async () => ({ error: 'no' }) }, + ); + + expect(listObservationsPage).toHaveBeenCalledWith(undefined, { + formType: 'demo', + includeDeleted: expected, + limit: 5000, + offset: 0, + }); + expect(postMessage).toHaveBeenCalledTimes(1); + }, + ); + it('routes replies to iframe matched by resolveReplyIframe when source is nested', async () => { const postPrimary = vi.fn(); const postNested = vi.fn(); diff --git a/desktop/src/lib/formPreviewBridge.ts b/desktop/src/lib/formPreviewBridge.ts index cd1bbe65b..950960f06 100644 --- a/desktop/src/lib/formPreviewBridge.ts +++ b/desktop/src/lib/formPreviewBridge.ts @@ -322,7 +322,7 @@ export function mapObservationToFormObservation( updatedAt: new Date(updated ?? r.lastSavedAt), syncedAt: new Date(synced ?? r.lastSavedAt), isDraft: false, - deleted: r.extras?.deleted === true, + deleted: r.deleted, formType: r.formType ?? '', formVersion: r.extras?.formVersion ?? '', data, @@ -491,15 +491,12 @@ export async function handleFormPreviewBridgeMessage( const includeDeleted = Boolean(data.includeDeleted); const page = await tauriClient.listObservationsPage(undefined, { formType, + includeDeleted, limit: 5000, offset: 0, }); - let rows = page.rows; - if (!includeDeleted) { - rows = rows.filter(r => r.extras?.deleted !== true); - } reply('getObservations', { - result: rows.map(mapObservationToFormObservation), + result: page.rows.map(mapObservationToFormObservation), }); return; } diff --git a/desktop/src/lib/tauriClient.ts b/desktop/src/lib/tauriClient.ts index 06819180e..8a7e98def 100644 --- a/desktop/src/lib/tauriClient.ts +++ b/desktop/src/lib/tauriClient.ts @@ -89,6 +89,7 @@ export const tauriClient = { query?: string, options?: { formType?: string | null; + includeDeleted?: boolean; limit?: number; offset?: number; }, @@ -96,6 +97,7 @@ export const tauriClient = { invokeSafe('list_observations_page', { query, formType: options?.formType ?? null, + includeDeleted: options?.includeDeleted, limit: options?.limit, offset: options?.offset, }), diff --git a/desktop/src/pages/ObservationsPage.tsx b/desktop/src/pages/ObservationsPage.tsx index 55885a4e9..b7124b944 100644 --- a/desktop/src/pages/ObservationsPage.tsx +++ b/desktop/src/pages/ObservationsPage.tsx @@ -39,7 +39,7 @@ function toPayloadText(value: unknown) { } function statusClass(obs: ObservationRecord) { - if (obs.extras?.deleted) return 'danger'; + if (obs.deleted) return 'danger'; if (obs.syncStatus === 'conflict') return 'danger'; if (obs.dirty) return 'warn'; return 'ok'; @@ -102,7 +102,7 @@ function draftFromRecord(record: ObservationRecord): ObservationEditorDraft { updatedAt: record.updatedAt ?? '', formVersion: x?.formVersion ?? DEFAULT_OBSERVATION_FORM_VERSION, createdAt: x?.createdAt ?? record.updatedAt ?? '', - deleted: x?.deleted ?? false, + deleted: record.deleted, syncedAt: x?.syncedAt ?? '', geoText: x?.geolocation != null && typeof x.geolocation === 'object' @@ -266,7 +266,7 @@ export function ObservationsPage() { } else if (filter === 'recent') { list = list.filter(isRecentlyModified); } else if (filter === 'deleted') { - list = list.filter(o => o.extras?.deleted); + list = list.filter(o => o.deleted); } return list; }, [observations, filter]); @@ -880,7 +880,7 @@ export function ObservationsPage() { title={syncPillLabel(item)}> {syncPillLabel(item)} - {item.extras?.deleted ? ( + {item.deleted ? ( Deleted ) : null} {item.id} diff --git a/desktop/src/services/synk/GeneratedSyncGateway.ts b/desktop/src/services/synk/GeneratedSyncGateway.ts index 3b86ff103..16f4042d0 100644 --- a/desktop/src/services/synk/GeneratedSyncGateway.ts +++ b/desktop/src/services/synk/GeneratedSyncGateway.ts @@ -90,7 +90,7 @@ function mapObservationToOpenApi(observation: ObservationRecord): Observation { data: payloadObject, created_at: createdAt, updated_at: updatedAt, - deleted: x?.deleted ?? false, + deleted: observation.deleted, synced_at: x?.syncedAt ? parseMaybeDate(x.syncedAt) : null, geolocation: geo, author: x?.author ?? null, diff --git a/desktop/src/types/domain.ts b/desktop/src/types/domain.ts index 4c59202de..c22d6d5dc 100644 --- a/desktop/src/types/domain.ts +++ b/desktop/src/types/domain.ts @@ -37,6 +37,8 @@ export interface ObservationRecord { hasConflictCopy: boolean; lastSavedAt: string; lastPushedAt?: string | null; + /** Authoritative local tombstone flag; mirrored in `extras.deleted` for envelope compatibility. */ + deleted: boolean; extras?: ObservationExtras | null; } diff --git a/packages/observation-query/src/compiler.test.ts b/packages/observation-query/src/compiler.test.ts index bd4532321..c4a420904 100644 --- a/packages/observation-query/src/compiler.test.ts +++ b/packages/observation-query/src/compiler.test.ts @@ -88,6 +88,54 @@ describe('ObservationQueryCompiler form_type wildcard', () => { }); }); +describe('ObservationQueryCompiler deleted column', () => { + it('uses the first-class Desktop deleted column for live-row filtering', () => { + const result = compileObservationQuery({ + dialect: 'desktop', + jsonColumn: 'payload', + indexKeys: new Set(), + formType: 'household', + includeDeleted: false, + }); + + expect('sql' in result).toBe(true); + if (!('sql' in result)) return; + expect(result.sql).toContain('o.deleted = 0'); + expect(result.sql).not.toContain('observation_extras'); + }); + + it('omits the deleted predicate when Desktop queries include deleted rows', () => { + const result = compileObservationQuery({ + dialect: 'desktop', + jsonColumn: 'payload', + indexKeys: new Set(), + formType: 'household', + includeDeleted: true, + }); + + expect('sql' in result).toBe(true); + if (!('sql' in result)) return; + expect(result.sql).not.toContain('o.deleted = 0'); + expect(result.sql).not.toContain('observation_extras'); + }); + + it('uses the first-class Desktop deleted column in metadata filters', () => { + const result = compileObservationQuery({ + dialect: 'desktop', + jsonColumn: 'payload', + indexKeys: new Set(), + includeDeleted: true, + filter: { field: 'deleted', op: 'eq', value: true }, + }); + + expect('sql' in result).toBe(true); + if (!('sql' in result)) return; + expect(result.sql).toContain('o.deleted = ?'); + expect(result.sql).not.toContain('observation_extras'); + expect(result.params).toEqual([true]); + }); +}); + describe('ObservationQueryCompiler fixtures', () => { const fixtures = loadFixtures(); diff --git a/packages/observation-query/src/compiler.ts b/packages/observation-query/src/compiler.ts index 0aebea4eb..0f50ed9a9 100644 --- a/packages/observation-query/src/compiler.ts +++ b/packages/observation-query/src/compiler.ts @@ -42,10 +42,7 @@ function metaColumnSql( return dialect === 'formulus' ? `${alias}.observation_id` : `${alias}.id`; } if (field === 'deleted') { - if (dialect === 'formulus') { - return `${alias}.deleted`; - } - return `COALESCE(json_extract(${alias}.observation_extras, '$.deleted'), 0)`; + return `${alias}.deleted`; } return `${alias}.${field}`; } @@ -301,13 +298,7 @@ export function compileObservationQuery( } if (!options.includeDeleted) { - if (dialect === 'formulus') { - whereParts.push(`${alias}.deleted = 0`); - } else { - whereParts.push( - `COALESCE(json_extract(${alias}.observation_extras, '$.deleted'), 0) = 0`, - ); - } + whereParts.push(`${alias}.deleted = 0`); } if (options.filter) {