From 487485bf1ea379ed1d87f3461a4129f204dcbf18 Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Tue, 8 Sep 2026 12:12:43 +0200 Subject: [PATCH 1/7] Add direct table option --- crates/core/src/schema/common.rs | 19 ++++++++ crates/core/src/schema/inspection.rs | 12 +++-- crates/core/src/schema/management.rs | 47 +++++++++++-------- crates/core/src/schema/raw_table.rs | 43 ++++++++++++++---- crates/core/src/schema/table_info.rs | 61 ++++++++++++++++++++++++- crates/core/src/views.rs | 13 ++++++ dart/test/schema_test.dart | 68 ++++++++++++++++++++++++++++ 7 files changed, 231 insertions(+), 32 deletions(-) diff --git a/crates/core/src/schema/common.rs b/crates/core/src/schema/common.rs index f133d44..2e47d25 100644 --- a/crates/core/src/schema/common.rs +++ b/crates/core/src/schema/common.rs @@ -18,6 +18,25 @@ pub enum SchemaTable<'a> { } impl<'a> SchemaTable<'a> { + /// The type name used for the table when referenced in `ps_crud`, `ps_oplog` and other tables. + pub fn name(&self) -> &str { + match self { + SchemaTable::Json(table) => &table.name, + SchemaTable::Raw { + definition, + schema: _, + } => &definition.name, + } + } + + pub fn data_column(&self) -> Option<&'static str> { + if let SchemaTable::Json(table) = self { + Some(table.data_column_name()) + } else { + None + } + } + pub fn common_options(&self) -> &CommonTableOptions { match self { Self::Json(table) => &table.options, diff --git a/crates/core/src/schema/inspection.rs b/crates/core/src/schema/inspection.rs index 09b72f7..1cc39fb 100644 --- a/crates/core/src/schema/inspection.rs +++ b/crates/core/src/schema/inspection.rs @@ -12,7 +12,9 @@ pub struct ExistingView { /// The name of the view itself. pub name: String, /// SQL contents of the `CREATE VIEW` statement. - pub sql: String, + /// + /// This is not set for as_raw_table tables, which don't have a view. + pub sql: Option, /// SQL contents of all triggers implementing deletes by forwarding to /// `ps_data` and `ps_crud`. pub delete_trigger_sql: String, @@ -52,7 +54,7 @@ SELECT results.push(ExistingView { name, - sql, + sql: Some(sql), delete_trigger_sql: delete, insert_trigger_sql: insert, update_trigger_sql: update, @@ -69,8 +71,10 @@ SELECT } pub fn create(&self, db: Database) -> Result<()> { - Self::drop_by_name(db, &self.name)?; - db.exec_safe_str(&self.sql)?; + if let Some(create_view) = &self.sql { + Self::drop_by_name(db, &self.name)?; + db.exec_safe_str(create_view)?; + } db.exec_safe_str(&self.delete_trigger_sql)?; db.exec_safe_str(&self.insert_trigger_sql)?; db.exec_safe_str(&self.update_trigger_sql)?; diff --git a/crates/core/src/schema/management.rs b/crates/core/src/schema/management.rs index 90f5d8c..20e742a 100644 --- a/crates/core/src/schema/management.rs +++ b/crates/core/src/schema/management.rs @@ -54,28 +54,33 @@ fn update_tables(db: Database, schema: &Schema) -> Result<()> { } // New table. - let quoted_internal_name = SqlBuffer::quote_identifier(&table.internal_name()); + let data_column = table.data_column_name(); + let mut create_table = SqlBuffer::default(); + + create_table.push_str("CREATE TABLE "); + table.write_name(&mut create_table); + _ = write!( + &mut create_table, + "(id TEXT PRIMARY KEY NOT NULL, {data_column} TEXT" + ); + + if table.direct { + for column in &table.columns { + create_table.push_char(','); + let _ = create_table.identifier().write_str(&column.name); + let _ = write!(&mut create_table, " {}", column.type_name); + } - db.exec_safe_str(&format!( - "CREATE TABLE {:}(id TEXT PRIMARY KEY NOT NULL, data TEXT)", - quoted_internal_name - ))?; + create_table.push_str(") STRICT /* ps-managed */;"); + } else { + create_table.push_str(");"); + } + + db.exec_safe_str(&create_table.sql)?; if !table.local_only() { // MOVE data if any - db.exec_text( - &format!( - "INSERT INTO {:}(id, data) - SELECT id, data - FROM ps_untyped - WHERE type = ?", - quoted_internal_name - ), - &table.name, - )?; - - // language=SQLite - db.exec_text("DELETE FROM ps_untyped WHERE type = ?", &table.name)?; + table.move_from_ps_untyped(db)?; } } @@ -217,7 +222,11 @@ fn update_views(db: Database, schema: &Schema) -> Result<()> { }; for table in &schema.tables { - let view_sql = powersync_view_sql(table); + let view_sql = if table.direct { + None + } else { + Some(powersync_view_sql(table)) + }; let delete_trigger_sql = powersync_trigger_delete_sql(table)?; let insert_trigger_sql = powersync_trigger_insert_sql(table)?; let update_trigger_sql = powersync_trigger_update_sql(table)?; diff --git a/crates/core/src/schema/raw_table.rs b/crates/core/src/schema/raw_table.rs index ad146aa..9d9cd3a 100644 --- a/crates/core/src/schema/raw_table.rs +++ b/crates/core/src/schema/raw_table.rs @@ -232,6 +232,22 @@ pub fn generate_raw_table_trigger( schema: &resolved_table, }; + generate_schema_table_trigger( + local_table_name, + as_schema_table, + synced_columns.as_ref(), + trigger_name, + write, + ) +} + +pub fn generate_schema_table_trigger( + local_table_name: &str, + table: SchemaTable, + synced_columns: Option<&ColumnFilter>, + trigger_name: &str, + write: WriteType, +) -> Result { let mut buffer = SqlBuffer::new(); buffer.create_trigger("", trigger_name); buffer.trigger_after(write, local_table_name); @@ -242,7 +258,7 @@ pub fn generate_raw_table_trigger( buffer.push_str(" AND\n("); // If we have a filter for synced columns (instead of syncing all of them), we want to add // additional WHEN clauses to enesure the trigger runs for updates on those columns only. - for (i, name) in as_schema_table.column_names().enumerate() { + for (i, name) in table.column_names().enumerate() { if i != 0 { buffer.push_str(" OR "); } @@ -258,14 +274,14 @@ pub fn generate_raw_table_trigger( buffer.push_str(" BEGIN\n"); - if table.schema.options.flags.insert_only() { + if table.common_options().flags.insert_only() { if write != WriteType::Insert { // Prevent illegal writes to a table marked as insert-only by raising errors here. buffer.push_str("SELECT RAISE(FAIL, 'Unexpected update on insert-only table');\n"); } else { // Insert-only tables use manual CRUD writes so they don't block incoming data. - let fragment = table_columns_to_json_object("NEW", &as_schema_table)?; - buffer.powersync_crud_manual_put(&table.name, &fragment); + let fragment = table_columns_to_json_object("NEW", &table)?; + buffer.powersync_crud_manual_put(table.name(), &fragment); } } else { if write == WriteType::Update { @@ -273,9 +289,9 @@ pub fn generate_raw_table_trigger( buffer.check_id_not_changed(); } - let json_fragment_new = table_columns_to_json_object("NEW", &as_schema_table)?; + let json_fragment_new = table_columns_to_json_object("NEW", &table)?; let json_fragment_old = if write == WriteType::Update { - Some(table_columns_to_json_object("OLD", &as_schema_table)?) + Some(table_columns_to_json_object("OLD", &table)?) } else { None }; @@ -294,15 +310,26 @@ pub fn generate_raw_table_trigger( write!(f, ", {json_fragment_new}))") }); + if write == WriteType::Update + && let Some(data_column) = table.data_column() + { + // If the table has a __data column storing the full JSON row, we also need to update + // that. + let _ = write!( + &mut buffer, + "UPDATE {local_table_name} SET {data_column} = {json_fragment_new} WHERE id = NEW.id;\n" + ); + } + buffer.insert_into_powersync_crud(InsertIntoCrud { op: write, - table: &as_schema_table, + table: &table, id_expr: if write == WriteType::Delete { "OLD.id" } else { "NEW.id" }, - type_name: &table.name, + type_name: table.name(), data: match write { // There is no data for deleted rows. WriteType::Delete => None, diff --git a/crates/core/src/schema/table_info.rs b/crates/core/src/schema/table_info.rs index fc46b9d..67bd96d 100644 --- a/crates/core/src/schema/table_info.rs +++ b/crates/core/src/schema/table_info.rs @@ -1,3 +1,5 @@ +use core::fmt::Write; + use alloc::rc::Rc; use alloc::string::ToString; use alloc::vec; @@ -5,7 +7,10 @@ use alloc::{collections::btree_set::BTreeSet, format, string::String, vec::Vec}; use serde::{Deserialize, de::Visitor}; use crate::error::PowerSyncError; -use crate::schema::ColumnFilter; +use crate::schema::raw_table::generate_schema_table_trigger; +use crate::schema::{ColumnFilter, SchemaTable}; +use crate::utils::database::Database; +use crate::utils::{SqlBuffer, WriteType}; #[derive(Deserialize)] pub struct Table { @@ -17,6 +22,7 @@ pub struct Table { pub indexes: Vec, #[serde(flatten)] pub options: CommonTableOptions, + pub direct: bool, } /// Options shared between regular and raw tables. @@ -78,6 +84,59 @@ impl Table { format!("ps_data__{:}", self.name) } } + + pub fn move_from_ps_untyped(&self, db: Database) -> Result<(), PowerSyncError> { + let mut stmt = SqlBuffer::default(); + let direct = self.direct; + + stmt.push_str("INSERT INTO "); + self.write_name(&mut stmt); + let _ = write!(&mut stmt, "(id, {}", self.data_column_name()); + + if direct { + for column in &self.columns { + stmt.push_char(','); + let _ = stmt.identifier().write_str(&column.name); + } + } + + stmt.push_str(") SELECT id, data"); + if direct { + for column in &self.columns { + stmt.push_char(','); + stmt.json_extract_and_cast("data", &column.name, &column.type_name); + } + } + + stmt.push_str(" FROM ps_untyped WHERE type = ?"); + + db.exec_text(&stmt.sql, &self.name)?; + db.exec_text("DELETE FROM ps_untyped WHERE type = ?", &self.name) + } + + pub fn write_name(&self, buffer: &mut SqlBuffer) { + if self.direct { + // Direct tables don't have views, so use the name of the table directly. + let _ = buffer.identifier().write_str(&self.name); + } else { + buffer.quote_internal_name(&self.name, self.local_only()); + } + } + + pub fn data_column_name(&self) -> &'static str { + if self.direct { "__data" } else { "data" } + } + + pub fn generate_direct_trigger(&self, write: WriteType) -> Result { + debug_assert!(self.direct); + generate_schema_table_trigger( + &self.name, + SchemaTable::Json(self), + None, + &format!("{}_trigger_{}", self.name, write), + write, + ) + } } impl RawTable { diff --git a/crates/core/src/views.rs b/crates/core/src/views.rs index a2c1b51..f5e7089 100644 --- a/crates/core/src/views.rs +++ b/crates/core/src/views.rs @@ -59,6 +59,10 @@ pub fn powersync_view_sql(table_info: &Table) -> String { } pub fn powersync_trigger_delete_sql(table_info: &Table) -> Result { + if table_info.direct { + return table_info.generate_direct_trigger(WriteType::Delete); + } + if table_info.options.flags.insert_only() { // Insert-only tables have no DELETE triggers return Ok(String::new()); @@ -117,6 +121,10 @@ pub fn powersync_trigger_delete_sql(table_info: &Table) -> Result { } pub fn powersync_trigger_insert_sql(table_info: &Table) -> Result { + if table_info.direct { + return table_info.generate_direct_trigger(WriteType::Insert); + } + let name = &table_info.name; let view_name = table_info.view_name(); let local_only = table_info.options.flags.local_only(); @@ -168,6 +176,10 @@ pub fn powersync_trigger_insert_sql(table_info: &Table) -> Result { } pub fn powersync_trigger_update_sql(table_info: &Table) -> Result { + if table_info.direct { + return table_info.generate_direct_trigger(WriteType::Update); + } + if table_info.options.flags.insert_only() { // Insert-only tables have no UPDATE triggers return Ok(String::new()); @@ -359,6 +371,7 @@ mod test { ], indexes: vec![], options: Default::default(), + direct: false, }; } diff --git a/dart/test/schema_test.dart b/dart/test/schema_test.dart index dfbd4ff..d925a4a 100644 --- a/dart/test/schema_test.dart +++ b/dart/test/schema_test.dart @@ -322,6 +322,74 @@ END''', test('#$i', () => testCase.testWith(db)); } }); + + group('direct tables', () { + final table = { + 'name': 'users', + 'columns': [ + {'name': 'name', 'type': 'text'} + ], + 'direct': true, + }; + + test('create', () { + db.executeInTx('SELECT powersync_replace_schema(?)', [ + json.encode({'tables': []}) + ]); + db.execute('INSERT INTO ps_untyped (type, id, data) VALUES (?, ?, ?)', [ + 'users', + 'user-id', + json.encode({'name': 'Name', 'other': 3}) + ]); + db.executeInTx('SELECT powersync_replace_schema(?)', [ + json.encode({ + 'tables': [table] + }) + ]); + + expect(db.select('SELECT * FROM users'), [ + { + 'id': 'user-id', + 'name': 'Name', + '__data': '{"name":"Name","other":3}' + }, + ]); + + final createTable = db.select( + 'SELECT sql FROM sqlite_schema WHERE type = ? AND tbl_name = ?', + ['table', 'users'], + )[0].columnAt(0); + expect( + createTable, + 'CREATE TABLE "users"(id TEXT PRIMARY KEY NOT NULL, __data TEXT,"name" text) STRICT /* ps-managed */', + ); + + final triggers = db + .select( + 'SELECT sql FROM sqlite_schema WHERE type = ? AND tbl_name = ? ORDER BY name', + ['trigger', 'users'], + ) + .map((r) => r['sql']) + .toList(); + + expect(triggers, [ + r''' +CREATE TRIGGER "users_trigger_DELETE" AFTER DELETE ON "users" FOR EACH ROW WHEN NOT powersync_in_sync_operation() BEGIN +INSERT INTO powersync_crud(op,id,type) VALUES ('DELETE', OLD.id, 'users'); +END''', + r''' +CREATE TRIGGER "users_trigger_INSERT" AFTER INSERT ON "users" FOR EACH ROW WHEN NOT powersync_in_sync_operation() BEGIN +INSERT INTO powersync_crud(op,id,type,data) VALUES ('PUT', NEW.id, 'users', json(powersync_diff('{}', json_object('name', powersync_strip_subtype(NEW."name"))))); +END''', + r''' +CREATE TRIGGER "users_trigger_UPDATE" AFTER UPDATE ON "users" FOR EACH ROW WHEN NOT powersync_in_sync_operation() BEGIN +SELECT CASE WHEN (OLD.id != NEW.id) THEN RAISE (FAIL, 'Cannot update id') END; +UPDATE users SET __data = json_object('name', powersync_strip_subtype(NEW."name")) WHERE id = NEW.id; +INSERT INTO powersync_crud(op,id,type,data,options) VALUES ('PATCH', NEW.id, 'users', json(powersync_diff(json_object('name', powersync_strip_subtype(OLD."name")), json_object('name', powersync_strip_subtype(NEW."name")))), 0); +END''' + ]); + }); + }); }); } From ebc95f57054706566016e9beae1c1115a0f983e6 Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Tue, 8 Sep 2026 13:46:47 +0200 Subject: [PATCH 2/7] Prepare sync local support --- crates/core/src/schema/common.rs | 116 +++++++++++++++++++++++- crates/core/src/schema/raw_table.rs | 126 +++------------------------ crates/core/src/schema/table_info.rs | 2 + crates/core/src/sync/sync_local.rs | 22 ++++- 4 files changed, 145 insertions(+), 121 deletions(-) diff --git a/crates/core/src/schema/common.rs b/crates/core/src/schema/common.rs index 2e47d25..846ce6e 100644 --- a/crates/core/src/schema/common.rs +++ b/crates/core/src/schema/common.rs @@ -1,10 +1,18 @@ -use core::slice; +use core::{fmt::Write, slice}; -use alloc::{string::String, vec::Vec}; +use alloc::{ + string::{String, ToString}, + vec, + vec::Vec, +}; use serde::Deserialize; -use crate::schema::{ - Column, CommonTableOptions, RawTable, Table, raw_table::InferredTableStructure, +use crate::{ + schema::{ + Column, CommonTableOptions, PendingStatement, PendingStatementValue, RawTable, Table, + raw_table::InferredTableStructure, + }, + utils::SqlBuffer, }; /// Utility to wrap both PowerSync-managed JSON tables and raw tables (with their schema snapshot @@ -57,6 +65,57 @@ impl<'a> SchemaTable<'a> { } => SchemaTableColumnIterator::Raw(schema.columns.iter()), } } + + /// Generates a statement of the form `INSERT INTO $tbl ($cols) VALUES (?, ...) ON CONFLICT (id) + /// DO UPDATE SET ...` for the sync client. + pub fn infer_put_stmt(&self, table_name: &str) -> PendingStatement { + let mut buffer = SqlBuffer::new(); + let mut params = vec![]; + + buffer.push_str("INSERT INTO "); + let _ = buffer.identifier().write_str(table_name); + buffer.push_str(" (id"); + for column in self.column_names() { + buffer.comma(); + let _ = buffer.identifier().write_str(column); + } + buffer.push_str(") VALUES (?1"); + params.push(PendingStatementValue::Id); + for (i, column) in self.column_names().enumerate() { + buffer.comma(); + let _ = write!(&mut buffer, "?{}", i + 2); + params.push(PendingStatementValue::Column(column.to_string())); + } + buffer.push_str(") ON CONFLICT (id) DO UPDATE SET "); + let mut do_update = buffer.comma_separated(); + // Generated an "x" = ? for all synced columns to update them without affecting local-only + // columns. + for (i, column) in self.column_names().enumerate() { + let entry = do_update.element(); + let _ = entry.identifier().write_str(column); + let _ = write!(entry, " = ?{}", i + 2); + } + + PendingStatement { + sql: buffer.sql, + params, + named_parameters_index: None, + } + } + + /// Generates a statement of the form `DELETE FROM $tbl WHERE id = ?` for the sync client. + pub fn infer_delete_stmt(&self, table_name: &str) -> PendingStatement { + let mut buffer = SqlBuffer::new(); + buffer.push_str("DELETE FROM "); + let _ = buffer.identifier().write_str(table_name); + buffer.push_str(" WHERE id = ?"); + + PendingStatement { + sql: buffer.sql, + params: vec![PendingStatementValue::Id], + named_parameters_index: None, + } + } } impl<'a> From<&'a Table> for SchemaTable<'a> { @@ -118,3 +177,52 @@ impl<'de> Deserialize<'de> for ColumnFilter { Ok(Self::from(Vec::::deserialize(deserializer)?)) } } +#[cfg(test)] +mod test { + use alloc::{string::ToString, vec}; + use core::assert_matches; + + use crate::schema::{ + PendingStatementValue, RawTable, SchemaTable, raw_table::InferredTableStructure, + table_info::RawTableSchema, + }; + + #[test] + fn infer_sync_statements() { + let raw_table = RawTable { + name: "users".to_string(), + schema: RawTableSchema::default(), + put: None, + delete: None, + clear: None, + }; + let structure = InferredTableStructure { + columns: vec!["foo".to_string(), "bar".to_string()], + }; + let schema_table = SchemaTable::Raw { + definition: &raw_table, + schema: &structure, + }; + + let put = schema_table.infer_put_stmt("tbl"); + assert_eq!( + put.sql, + r#"INSERT INTO "tbl" (id, "foo", "bar") VALUES (?1, ?2, ?3) ON CONFLICT (id) DO UPDATE SET "foo" = ?2, "bar" = ?3"# + ); + assert_eq!(put.params.len(), 3); + assert_matches!(put.params[0], PendingStatementValue::Id); + assert_matches!( + put.params[1], + PendingStatementValue::Column(ref name) if name == "foo" + ); + assert_matches!( + put.params[2], + PendingStatementValue::Column(ref name) if name == "bar" + ); + + let delete = schema_table.infer_delete_stmt("tbl"); + assert_eq!(delete.sql, r#"DELETE FROM "tbl" WHERE id = ?"#); + assert_eq!(delete.params.len(), 1); + assert_matches!(delete.params[0], PendingStatementValue::Id); + } +} diff --git a/crates/core/src/schema/raw_table.rs b/crates/core/src/schema/raw_table.rs index 9d9cd3a..878b23f 100644 --- a/crates/core/src/schema/raw_table.rs +++ b/crates/core/src/schema/raw_table.rs @@ -15,13 +15,12 @@ use powersync_sqlite_nostd::Destructor; use crate::{ error::{PowerSyncError, Result}, - schema::{ColumnFilter, PendingStatement, PendingStatementValue, RawTable, SchemaTable}, + schema::{ColumnFilter, PendingStatement, RawTable, SchemaTable}, utils::{InsertIntoCrud, SqlBuffer, WriteType, database::Database}, views::table_columns_to_json_object, }; pub struct InferredTableStructure { - pub name: String, pub columns: Vec, } @@ -59,61 +58,7 @@ impl InferredTableStructure { "Table {table_name} has no id column." ))) } else { - Ok(Self { - name: table_name.to_string(), - columns, - }) - } - } - - /// Generates a statement of the form `INSERT INTO $tbl ($cols) VALUES (?, ...) ON CONFLICT (id) - /// DO UPDATE SET ...` for the sync client. - pub fn infer_put_stmt(&self) -> PendingStatement { - let mut buffer = SqlBuffer::new(); - let mut params = vec![]; - - buffer.push_str("INSERT INTO "); - let _ = buffer.identifier().write_str(&self.name); - buffer.push_str(" (id"); - for column in &self.columns { - buffer.comma(); - let _ = buffer.identifier().write_str(column); - } - buffer.push_str(") VALUES (?1"); - params.push(PendingStatementValue::Id); - for (i, column) in self.columns.iter().enumerate() { - buffer.comma(); - let _ = write!(&mut buffer, "?{}", i + 2); - params.push(PendingStatementValue::Column(column.clone())); - } - buffer.push_str(") ON CONFLICT (id) DO UPDATE SET "); - let mut do_update = buffer.comma_separated(); - // Generated an "x" = ? for all synced columns to update them without affecting local-only - // columns. - for (i, column) in self.columns.iter().enumerate() { - let entry = do_update.element(); - let _ = entry.identifier().write_str(column); - let _ = write!(entry, " = ?{}", i + 2); - } - - PendingStatement { - sql: buffer.sql, - params, - named_parameters_index: None, - } - } - - /// Generates a statement of the form `DELETE FROM $tbl WHERE id = ?` for the sync client. - pub fn infer_delete_stmt(&self) -> PendingStatement { - let mut buffer = SqlBuffer::new(); - buffer.push_str("DELETE FROM "); - let _ = buffer.identifier().write_str(&self.name); - buffer.push_str(" WHERE id = ?"); - - PendingStatement { - sql: buffer.sql, - params: vec![PendingStatementValue::Id], - named_parameters_index: None, + Ok(Self { columns }) } } } @@ -141,7 +86,7 @@ impl InferredSchemaCache { schema_version: usize, tbl: &RawTable, ) -> Result> { - self.with_entry(db, schema_version, tbl, SchemaCacheEntry::put) + self.with_entry(db, schema_version, tbl, |entry| entry.put_stmt.clone()) } pub fn infer_delete_statement( @@ -150,7 +95,7 @@ impl InferredSchemaCache { schema_version: usize, tbl: &RawTable, ) -> Result> { - self.with_entry(db, schema_version, tbl, SchemaCacheEntry::delete) + self.with_entry(db, schema_version, tbl, |entry| entry.delete_stmt.clone()) } fn with_entry( @@ -179,9 +124,8 @@ impl InferredSchemaCache { pub struct SchemaCacheEntry { schema_version: usize, - structure: InferredTableStructure, - put_stmt: Option>, - delete_stmt: Option>, + pub put_stmt: Rc, + pub delete_stmt: Rc, } impl SchemaCacheEntry { @@ -192,26 +136,17 @@ impl SchemaCacheEntry { db, &table.schema.synced_columns, )?; + let schema_table = SchemaTable::Raw { + definition: table, + schema: &structure, + }; Ok(Self { schema_version, - structure, - put_stmt: None, - delete_stmt: None, + put_stmt: Rc::new(schema_table.infer_put_stmt(local_table_name)), + delete_stmt: Rc::new(schema_table.infer_delete_stmt(local_table_name)), }) } - - fn put(&mut self) -> Rc { - self.put_stmt - .get_or_insert_with(|| Rc::new(self.structure.infer_put_stmt())) - .clone() - } - - fn delete(&mut self) -> Rc { - self.delete_stmt - .get_or_insert_with(|| Rc::new(self.structure.infer_delete_stmt())) - .clone() - } } /// Generates a `CREATE TRIGGER` statement to capture writes on raw tables and to forward them to @@ -342,40 +277,3 @@ pub fn generate_schema_table_trigger( buffer.trigger_end(); Ok(buffer.sql) } - -#[cfg(test)] -mod test { - use alloc::{string::ToString, vec}; - use core::assert_matches; - - use crate::schema::{PendingStatementValue, raw_table::InferredTableStructure}; - - #[test] - fn infer_sync_statements() { - let structure = InferredTableStructure { - name: "tbl".to_string(), - columns: vec!["foo".to_string(), "bar".to_string()], - }; - - let put = structure.infer_put_stmt(); - assert_eq!( - put.sql, - r#"INSERT INTO "tbl" (id, "foo", "bar") VALUES (?1, ?2, ?3) ON CONFLICT (id) DO UPDATE SET "foo" = ?2, "bar" = ?3"# - ); - assert_eq!(put.params.len(), 3); - assert_matches!(put.params[0], PendingStatementValue::Id); - assert_matches!( - put.params[1], - PendingStatementValue::Column(ref name) if name == "foo" - ); - assert_matches!( - put.params[2], - PendingStatementValue::Column(ref name) if name == "bar" - ); - - let delete = structure.infer_delete_stmt(); - assert_eq!(delete.sql, r#"DELETE FROM "tbl" WHERE id = ?"#); - assert_eq!(delete.params.len(), 1); - assert_matches!(delete.params[0], PendingStatementValue::Id); - } -} diff --git a/crates/core/src/schema/table_info.rs b/crates/core/src/schema/table_info.rs index 67bd96d..858d330 100644 --- a/crates/core/src/schema/table_info.rs +++ b/crates/core/src/schema/table_info.rs @@ -78,6 +78,8 @@ impl Table { } pub fn internal_name(&self) -> String { + debug_assert!(!self.direct); + if self.local_only() { format!("ps_data_local__{:}", self.name) } else { diff --git a/crates/core/src/sync/sync_local.rs b/crates/core/src/sync/sync_local.rs index fb3c449..5ba78cd 100644 --- a/crates/core/src/sync/sync_local.rs +++ b/crates/core/src/sync/sync_local.rs @@ -11,7 +11,8 @@ use serde::ser::SerializeMap; use crate::error::{PowerSyncError, Result}; use crate::schema::inspection::ExistingTable; use crate::schema::{ - InferredSchemaCache, PendingStatement, PendingStatementValue, RawTable, Schema, + InferredSchemaCache, PendingStatement, PendingStatementValue, RawTable, Schema, SchemaTable, + Table, }; use crate::state::DatabaseState; use crate::sync::BucketPriority; @@ -124,7 +125,6 @@ WHERE target.key = '{TARGET_CHECKPOINT_REQUEST_ID_KEY}' "expected oplog data to be an object", ) })?; - let rest = stmt.render_rest_object(json_object)?; stmt.bind_for_put(id, data, Some(json_object), rest.as_ref())?; stmt.exec(type_name, id, Some(&data))?; @@ -340,6 +340,15 @@ impl<'a> ParsedDatabaseSchema<'a> { } fn add_from_schema(&mut self, schema: &'a Schema) { + for regular in &schema.tables { + if regular.direct { + self.tables.insert( + regular.name.clone(), + ParsedSchemaTable::new(TableDefinition::Direct(regular)), + ); + } + } + for raw in &schema.raw_tables { self.tables.insert( raw.name.clone(), @@ -420,6 +429,9 @@ impl<'a> ParsedSchemaTable<'a> { named_parameters_index: None, }) } + TableDefinition::Direct(table) => { + Rc::new(SchemaTable::Json(table).infer_put_stmt(&table.name)) + } }) }) } @@ -448,6 +460,9 @@ impl<'a> ParsedSchemaTable<'a> { named_parameters_index: None, }) } + TableDefinition::Direct(table) => { + Rc::new(SchemaTable::Json(table).infer_delete_stmt(&table.name)) + } }) }) } @@ -456,12 +471,13 @@ impl<'a> ParsedSchemaTable<'a> { enum TableDefinition<'a> { Raw(&'a RawTable), JsonView { local_table: String }, + Direct(&'a Table), } struct PreparedPendingStatement { stmt: Statement, - definition: Rc, needs_parsed_json: bool, + definition: Rc, } impl PreparedPendingStatement { From d22f333b1241aa48f40f23f95e0c494342b5726b Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Tue, 8 Sep 2026 18:03:00 +0200 Subject: [PATCH 3/7] Simple sync test --- crates/core/src/schema/common.rs | 23 ++++++++++-- crates/core/src/schema/table_info.rs | 3 +- dart/test/sync_test.dart | 55 ++++++++++++++++++++++++++++ 3 files changed, 76 insertions(+), 5 deletions(-) diff --git a/crates/core/src/schema/common.rs b/crates/core/src/schema/common.rs index 846ce6e..7165859 100644 --- a/crates/core/src/schema/common.rs +++ b/crates/core/src/schema/common.rs @@ -71,29 +71,46 @@ impl<'a> SchemaTable<'a> { pub fn infer_put_stmt(&self, table_name: &str) -> PendingStatement { let mut buffer = SqlBuffer::new(); let mut params = vec![]; + let data_column = self.data_column(); buffer.push_str("INSERT INTO "); let _ = buffer.identifier().write_str(table_name); buffer.push_str(" (id"); + if let Some(data_column) = data_column { + let _ = write!(&mut buffer, ", {data_column}"); + } + for column in self.column_names() { buffer.comma(); let _ = buffer.identifier().write_str(column); } buffer.push_str(") VALUES (?1"); params.push(PendingStatementValue::Id); + if data_column.is_some() { + params.push(PendingStatementValue::Row); + buffer.push_str(", ?2"); + } + + let data_start_index = if data_column.is_some() { 3 } else { 2 }; for (i, column) in self.column_names().enumerate() { buffer.comma(); - let _ = write!(&mut buffer, "?{}", i + 2); + let _ = write!(&mut buffer, "?{}", i + data_start_index); params.push(PendingStatementValue::Column(column.to_string())); } buffer.push_str(") ON CONFLICT (id) DO UPDATE SET "); let mut do_update = buffer.comma_separated(); - // Generated an "x" = ? for all synced columns to update them without affecting local-only + + if let Some(data_column) = data_column { + let entry = do_update.element(); + let _ = write!(entry, "{data_column} = ?2"); + } + + // Generate an "x" = ? for all synced columns to update them without affecting local-only // columns. for (i, column) in self.column_names().enumerate() { let entry = do_update.element(); let _ = entry.identifier().write_str(column); - let _ = write!(entry, " = ?{}", i + 2); + let _ = write!(entry, " = ?{}", i + data_start_index); } PendingStatement { diff --git a/crates/core/src/schema/table_info.rs b/crates/core/src/schema/table_info.rs index 858d330..879dab3 100644 --- a/crates/core/src/schema/table_info.rs +++ b/crates/core/src/schema/table_info.rs @@ -22,6 +22,7 @@ pub struct Table { pub indexes: Vec, #[serde(flatten)] pub options: CommonTableOptions, + #[serde(default)] pub direct: bool, } @@ -78,8 +79,6 @@ impl Table { } pub fn internal_name(&self) -> String { - debug_assert!(!self.direct); - if self.local_only() { format!("ps_data_local__{:}", self.name) } else { diff --git a/dart/test/sync_test.dart b/dart/test/sync_test.dart index b7c969b..4aaf199 100644 --- a/dart/test/sync_test.dart +++ b/dart/test/sync_test.dart @@ -2180,6 +2180,61 @@ CREATE TRIGGER users_ref_delete }); }); + group('direct tables', () { + test('smoke test', () { + final schema = { + 'tables': [ + { + 'name': 'users', + 'columns': [ + {'name': 'name', 'type': 'text'} + ], + 'direct': true, + } + ] + }; + + db.executeInTx( + 'SELECT powersync_replace_schema(?)', [json.encode(schema)]); + invokeControl('start', json.encode({'schema': schema})); + + // Insert + pushCheckpoint(buckets: [bucketDescription('a')]); + pushSyncData( + 'a', + '1', + 'my_user', + 'PUT', + {'name': 'First user'}, + objectType: 'users', + ); + pushCheckpointComplete(); + + final users = db.select('SELECT * FROM users;'); + expect(users, [ + { + 'id': 'my_user', + 'name': 'First user', + '__data': '{"name":"First user"}' + } + ]); + + // Delete + pushCheckpoint(buckets: [bucketDescription('a')]); + pushSyncData( + 'a', + '1', + 'my_user', + 'REMOVE', + null, + objectType: 'users', + ); + pushCheckpointComplete(); + + expect(db.select('SELECT * FROM users'), isEmpty); + }); + }); + test('can close database while iteration is active', () { // The sync client caches prepared statements, we need to ensure those are // freed when we close the connection since SQLite would keep files open From 80b2c57cb138e7f6da7f490aebdd4a86013d25c2 Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Wed, 9 Sep 2026 13:18:56 +0200 Subject: [PATCH 4/7] Allow removing tables --- crates/core/src/schema/common.rs | 43 +++++++++++++--------------- crates/core/src/schema/inspection.rs | 34 +++++++++++++++------- crates/core/src/schema/management.rs | 13 ++++----- crates/core/src/schema/raw_table.rs | 21 +++++++------- crates/core/src/schema/table_info.rs | 6 +++- crates/core/src/sync/sync_local.rs | 2 +- dart/test/schema_test.dart | 19 +++++++++++- 7 files changed, 84 insertions(+), 54 deletions(-) diff --git a/crates/core/src/schema/common.rs b/crates/core/src/schema/common.rs index 7165859..3600b22 100644 --- a/crates/core/src/schema/common.rs +++ b/crates/core/src/schema/common.rs @@ -1,4 +1,4 @@ -use core::{fmt::Write, slice}; +use core::fmt::Write; use alloc::{ string::{String, ToString}, @@ -55,17 +55,21 @@ impl<'a> SchemaTable<'a> { } } - /// Iterates over defined column names in this table (not including the `id` column). - pub fn column_names(&self) -> impl Iterator { + pub fn columns(&self) -> &'a [Column] { match self { - Self::Json(table) => SchemaTableColumnIterator::Json(table.columns.iter()), + Self::Json(table) => &table.columns, Self::Raw { definition: _, schema, - } => SchemaTableColumnIterator::Raw(schema.columns.iter()), + } => &schema.columns, } } + /// Iterates over defined column names in this table (not including the `id` column). + pub fn column_names(&self) -> impl Iterator { + self.columns().iter().map(|c| &*c.name) + } + /// Generates a statement of the form `INSERT INTO $tbl ($cols) VALUES (?, ...) ON CONFLICT (id) /// DO UPDATE SET ...` for the sync client. pub fn infer_put_stmt(&self, table_name: &str) -> PendingStatement { @@ -141,22 +145,6 @@ impl<'a> From<&'a Table> for SchemaTable<'a> { } } -enum SchemaTableColumnIterator<'a> { - Json(slice::Iter<'a, Column>), - Raw(slice::Iter<'a, String>), -} - -impl<'a> Iterator for SchemaTableColumnIterator<'a> { - type Item = &'a str; - - fn next(&mut self) -> Option { - Some(match self { - Self::Json(iter) => &iter.next()?.name, - Self::Raw(iter) => iter.next()?.as_ref(), - }) - } -} - #[derive(Default)] pub struct ColumnFilter { sorted_names: Vec, @@ -200,7 +188,7 @@ mod test { use core::assert_matches; use crate::schema::{ - PendingStatementValue, RawTable, SchemaTable, raw_table::InferredTableStructure, + Column, PendingStatementValue, RawTable, SchemaTable, raw_table::InferredTableStructure, table_info::RawTableSchema, }; @@ -214,7 +202,16 @@ mod test { clear: None, }; let structure = InferredTableStructure { - columns: vec!["foo".to_string(), "bar".to_string()], + columns: vec![ + Column { + name: "foo".to_string(), + type_name: "TEXT".to_string(), + }, + Column { + name: "bar".to_string(), + type_name: "TEXT".to_string(), + }, + ], }; let schema_table = SchemaTable::Raw { definition: &raw_table, diff --git a/crates/core/src/schema/inspection.rs b/crates/core/src/schema/inspection.rs index 1cc39fb..fc43d64 100644 --- a/crates/core/src/schema/inspection.rs +++ b/crates/core/src/schema/inspection.rs @@ -3,6 +3,7 @@ use alloc::{format, vec}; use alloc::{string::String, vec::Vec}; use crate::error::Result; +use crate::schema::raw_table::InferredTableStructure; use crate::utils::SqlBuffer; use crate::utils::database::Database; @@ -87,28 +88,39 @@ pub struct ExistingTable { pub name: String, pub internal_name: String, pub local_only: bool, + pub direct: Option, } impl ExistingTable { pub fn list(db: Database) -> Result> { let mut results = vec![]; - let stmt = db.prepare_v2( - " -SELECT name FROM sqlite_master WHERE type = 'table' AND name GLOB 'ps_data_*'; - ", - )?; + let stmt = db.prepare_v2("SELECT name, sql FROM sqlite_master WHERE type = 'table';")?; while stmt.step()? { let internal_name = stmt.column_text(0)?; - let Some((name, local_only)) = Self::external_name(internal_name) else { + let Ok(sql) = stmt.column_text(1) else { continue; }; - results.push(ExistingTable { - internal_name: internal_name.to_owned(), - name: name.to_owned(), - local_only: local_only, - }); + if let Some((name, local_only)) = Self::external_name(internal_name) { + results.push(ExistingTable { + internal_name: internal_name.to_owned(), + name: name.to_owned(), + local_only: local_only, + direct: None, + }); + } else if sql.contains("/* ps-managed */") { + results.push(ExistingTable { + internal_name: internal_name.to_owned(), + name: internal_name.to_owned(), + local_only: false, + direct: Some(InferredTableStructure::read_from_database( + internal_name, + db, + &None, + )?), + }); + } } Ok(results) diff --git a/crates/core/src/schema/management.rs b/crates/core/src/schema/management.rs index 20e742a..3c71717 100644 --- a/crates/core/src/schema/management.rs +++ b/crates/core/src/schema/management.rs @@ -16,7 +16,7 @@ use sqlite::{Connection, ResultCode, Value}; use crate::create_sqlite_text_fn; use crate::error::{PowerSyncError, Result}; use crate::schema::inspection::{ExistingTable, ExistingView}; -use crate::schema::table_info::Index; +use crate::schema::table_info::{Index, data_column_name}; use crate::state::DatabaseState; use crate::utils::database::Database; use crate::utils::{SqlBuffer, verify_in_transaction}; @@ -65,17 +65,15 @@ fn update_tables(db: Database, schema: &Schema) -> Result<()> { ); if table.direct { + create_table.push_str("/* ps-managed */"); + for column in &table.columns { create_table.push_char(','); let _ = create_table.identifier().write_str(&column.name); let _ = write!(&mut create_table, " {}", column.type_name); } - - create_table.push_str(") STRICT /* ps-managed */;"); - } else { - create_table.push_str(");"); } - + create_table.push_str(");"); db.exec_safe_str(&create_table.sql)?; if !table.local_only() { @@ -90,7 +88,8 @@ fn update_tables(db: Database, schema: &Schema) -> Result<()> { if !remaining.local_only { db.exec_text( &format!( - "INSERT INTO ps_untyped(type, id, data) SELECT ?, id, data FROM {:}", + "INSERT INTO ps_untyped(type, id, data) SELECT ?, id, {} FROM {:}", + data_column_name(remaining.direct.is_some()), SqlBuffer::quote_identifier(&remaining.internal_name) ), &remaining.name, diff --git a/crates/core/src/schema/raw_table.rs b/crates/core/src/schema/raw_table.rs index 878b23f..4168dfd 100644 --- a/crates/core/src/schema/raw_table.rs +++ b/crates/core/src/schema/raw_table.rs @@ -4,24 +4,20 @@ use core::{ }; use alloc::{ - collections::btree_map::BTreeMap, - format, - rc::Rc, - string::{String, ToString}, - vec, + borrow::ToOwned, collections::btree_map::BTreeMap, format, rc::Rc, string::String, vec, vec::Vec, }; use powersync_sqlite_nostd::Destructor; use crate::{ error::{PowerSyncError, Result}, - schema::{ColumnFilter, PendingStatement, RawTable, SchemaTable}, + schema::{Column, ColumnFilter, PendingStatement, RawTable, SchemaTable}, utils::{InsertIntoCrud, SqlBuffer, WriteType, database::Database}, views::table_columns_to_json_object, }; pub struct InferredTableStructure { - pub columns: Vec, + pub columns: Vec, } impl InferredTableStructure { @@ -30,7 +26,7 @@ impl InferredTableStructure { db: Database, synced_columns: &Option, ) -> Result { - let stmt = db.prepare_v2("select name from pragma_table_info(?)")?; + let stmt = db.prepare_v2("select name, type from pragma_table_info(?)")?; stmt.bind_text(1, table_name, Destructor::STATIC)?; let mut has_id_column = false; @@ -38,6 +34,8 @@ impl InferredTableStructure { while stmt.step()? { let name = stmt.column_text(0)?; + let column_type = stmt.column_text(1)?; + if name == "id" { has_id_column = true; } else if let Some(filter) = synced_columns @@ -45,7 +43,10 @@ impl InferredTableStructure { { // This column isn't part of the synced columns, skip. } else { - columns.push(name.to_string()); + columns.push(Column { + name: name.to_owned(), + type_name: column_type.to_owned(), + }); } } @@ -245,7 +246,7 @@ pub fn generate_schema_table_trigger( write!(f, ", {json_fragment_new}))") }); - if write == WriteType::Update + if write != WriteType::Delete && let Some(data_column) = table.data_column() { // If the table has a __data column storing the full JSON row, we also need to update diff --git a/crates/core/src/schema/table_info.rs b/crates/core/src/schema/table_info.rs index 879dab3..f4d3f0e 100644 --- a/crates/core/src/schema/table_info.rs +++ b/crates/core/src/schema/table_info.rs @@ -125,7 +125,7 @@ impl Table { } pub fn data_column_name(&self) -> &'static str { - if self.direct { "__data" } else { "data" } + data_column_name(self.direct) } pub fn generate_direct_trigger(&self, write: WriteType) -> Result { @@ -140,6 +140,10 @@ impl Table { } } +pub fn data_column_name(is_direct: bool) -> &'static str { + if is_direct { "__data" } else { "data" } +} + impl RawTable { pub fn require_table_name(&self) -> Result<&str, PowerSyncError> { let Some(local_table_name) = self.schema.table_name.as_ref() else { diff --git a/crates/core/src/sync/sync_local.rs b/crates/core/src/sync/sync_local.rs index 5ba78cd..fe68f72 100644 --- a/crates/core/src/sync/sync_local.rs +++ b/crates/core/src/sync/sync_local.rs @@ -360,7 +360,7 @@ impl<'a> ParsedDatabaseSchema<'a> { fn add_from_db(&mut self, db: Database) -> Result<()> { let tables = ExistingTable::list(db)?; for table in tables { - if !table.local_only { + if !table.local_only && !self.tables.contains_key(&table.name) { let visible_name = table.name; self.tables.insert( diff --git a/dart/test/schema_test.dart b/dart/test/schema_test.dart index d925a4a..a585acf 100644 --- a/dart/test/schema_test.dart +++ b/dart/test/schema_test.dart @@ -361,7 +361,7 @@ END''', )[0].columnAt(0); expect( createTable, - 'CREATE TABLE "users"(id TEXT PRIMARY KEY NOT NULL, __data TEXT,"name" text) STRICT /* ps-managed */', + 'CREATE TABLE "users"(id TEXT PRIMARY KEY NOT NULL, __data TEXT/* ps-managed */,"name" text)', ); final triggers = db @@ -389,6 +389,23 @@ INSERT INTO powersync_crud(op,id,type,data,options) VALUES ('PATCH', NEW.id, 'us END''' ]); }); + + test('remove from schema', () { + db.executeInTx('SELECT powersync_replace_schema(?)', [ + json.encode({ + 'tables': [table] + }) + ]); + db.execute( + 'INSERT INTO users (id, name) VALUES (?, ?)', ['id', 'name']); + db.executeInTx('SELECT powersync_replace_schema(?)', [ + json.encode({'tables': []}) + ]); + + expect(db.select('SELECT * FROM ps_untyped'), [ + {'type': 'users', 'id': 'id', 'data': '{"name":"name"}'} + ]); + }); }); }); } From 4350da58fdf260f4be0916fd3014beeece4dd9ae Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Wed, 9 Sep 2026 13:33:20 +0200 Subject: [PATCH 5/7] Support local-only direct tables --- crates/core/src/schema/inspection.rs | 6 +++- crates/core/src/schema/management.rs | 3 +- crates/core/src/schema/raw_table.rs | 47 +++++++++++++++++----------- crates/core/src/sync/sync_local.rs | 6 ++-- dart/test/schema_test.dart | 43 +++++++++++++++---------- dart/test/sync_test.dart | 33 ++++++++++++++++--- 6 files changed, 96 insertions(+), 42 deletions(-) diff --git a/crates/core/src/schema/inspection.rs b/crates/core/src/schema/inspection.rs index fc43d64..902320d 100644 --- a/crates/core/src/schema/inspection.rs +++ b/crates/core/src/schema/inspection.rs @@ -93,6 +93,10 @@ pub struct ExistingTable { impl ExistingTable { pub fn list(db: Database) -> Result> { + Self::list_filtered(db, false) + } + + pub fn list_filtered(db: Database, ignore_direct: bool) -> Result> { let mut results = vec![]; let stmt = db.prepare_v2("SELECT name, sql FROM sqlite_master WHERE type = 'table';")?; @@ -109,7 +113,7 @@ impl ExistingTable { local_only: local_only, direct: None, }); - } else if sql.contains("/* ps-managed */") { + } else if sql.contains("/* ps-managed */") && !ignore_direct { results.push(ExistingTable { internal_name: internal_name.to_owned(), name: internal_name.to_owned(), diff --git a/crates/core/src/schema/management.rs b/crates/core/src/schema/management.rs index 3c71717..24a841d 100644 --- a/crates/core/src/schema/management.rs +++ b/crates/core/src/schema/management.rs @@ -39,11 +39,12 @@ fn update_tables(db: Database, schema: &Schema) -> Result<()> { for table in &schema.tables { if let Some(existing) = existing_tables.remove(&*table.name) { - if existing.local_only != table.local_only() { + if !table.direct && existing.local_only != table.local_only() { // Migrating between local-only and synced tables. This works by deleting // existing and re-creating the table from scratch. We can re-create first and // delete the old table afterwards because they have a different name // (local-only tables have a ps_data_local prefix). + // Direct tables are the same whether they're local or not. // To delete the old existing table in the end. existing_tables.insert(&existing.name, existing); diff --git a/crates/core/src/schema/raw_table.rs b/crates/core/src/schema/raw_table.rs index 4168dfd..54dbdec 100644 --- a/crates/core/src/schema/raw_table.rs +++ b/crates/core/src/schema/raw_table.rs @@ -209,12 +209,14 @@ pub fn generate_schema_table_trigger( } buffer.push_str(" BEGIN\n"); + let flags = table.common_options().flags; + let mut has_stmt = false; - if table.common_options().flags.insert_only() { + if flags.insert_only() { if write != WriteType::Insert { // Prevent illegal writes to a table marked as insert-only by raising errors here. buffer.push_str("SELECT RAISE(FAIL, 'Unexpected update on insert-only table');\n"); - } else { + } else if !flags.local_only() { // Insert-only tables use manual CRUD writes so they don't block incoming data. let fragment = table_columns_to_json_object("NEW", &table)?; buffer.powersync_crud_manual_put(table.name(), &fragment); @@ -255,24 +257,33 @@ pub fn generate_schema_table_trigger( &mut buffer, "UPDATE {local_table_name} SET {data_column} = {json_fragment_new} WHERE id = NEW.id;\n" ); + + has_stmt = true; } - buffer.insert_into_powersync_crud(InsertIntoCrud { - op: write, - table: &table, - id_expr: if write == WriteType::Delete { - "OLD.id" - } else { - "NEW.id" - }, - type_name: table.name(), - data: match write { - // There is no data for deleted rows. - WriteType::Delete => None, - _ => Some(&write_data), - }, - metadata: None::<&'static str>, - })?; + if !flags.local_only() { + has_stmt = true; + buffer.insert_into_powersync_crud(InsertIntoCrud { + op: write, + table: &table, + id_expr: if write == WriteType::Delete { + "OLD.id" + } else { + "NEW.id" + }, + type_name: table.name(), + data: match write { + // There is no data for deleted rows. + WriteType::Delete => None, + _ => Some(&write_data), + }, + metadata: None::<&'static str>, + })?; + } + } + + if !has_stmt { + return Ok(Default::default()); } buffer.trigger_end(); diff --git a/crates/core/src/sync/sync_local.rs b/crates/core/src/sync/sync_local.rs index fe68f72..24f6b19 100644 --- a/crates/core/src/sync/sync_local.rs +++ b/crates/core/src/sync/sync_local.rs @@ -341,7 +341,7 @@ impl<'a> ParsedDatabaseSchema<'a> { fn add_from_schema(&mut self, schema: &'a Schema) { for regular in &schema.tables { - if regular.direct { + if regular.direct && !regular.local_only() { self.tables.insert( regular.name.clone(), ParsedSchemaTable::new(TableDefinition::Direct(regular)), @@ -358,7 +358,9 @@ impl<'a> ParsedDatabaseSchema<'a> { } fn add_from_db(&mut self, db: Database) -> Result<()> { - let tables = ExistingTable::list(db)?; + // Ignore direct tables here, we can rely on them being added via add_from_schema. + // TODO: Remove this function, SDKs should always pass the used schema when they connect. + let tables = ExistingTable::list_filtered(db, true)?; for table in tables { if !table.local_only && !self.tables.contains_key(&table.name) { let visible_name = table.name; diff --git a/dart/test/schema_test.dart b/dart/test/schema_test.dart index a585acf..f51106f 100644 --- a/dart/test/schema_test.dart +++ b/dart/test/schema_test.dart @@ -324,13 +324,20 @@ END''', }); group('direct tables', () { - final table = { - 'name': 'users', - 'columns': [ - {'name': 'name', 'type': 'text'} - ], - 'direct': true, - }; + Object schema({Map additionalOptions = const {}}) { + return { + 'tables': [ + { + 'name': 'users', + 'columns': [ + {'name': 'name', 'type': 'text'} + ], + 'direct': true, + ...additionalOptions, + } + ] + }; + } test('create', () { db.executeInTx('SELECT powersync_replace_schema(?)', [ @@ -341,11 +348,8 @@ END''', 'user-id', json.encode({'name': 'Name', 'other': 3}) ]); - db.executeInTx('SELECT powersync_replace_schema(?)', [ - json.encode({ - 'tables': [table] - }) - ]); + db.executeInTx( + 'SELECT powersync_replace_schema(?)', [json.encode(schema())]); expect(db.select('SELECT * FROM users'), [ { @@ -390,12 +394,19 @@ END''' ]); }); - test('remove from schema', () { + test('local-only', () { db.executeInTx('SELECT powersync_replace_schema(?)', [ - json.encode({ - 'tables': [table] - }) + json.encode(schema(additionalOptions: {'local_only': true})) ]); + + db.execute( + 'INSERT INTO users (id, name) VALUES (?, ?)', ['id', 'name']); + expect(db.select('SELECT * FROM ps_crud'), isEmpty); + }); + + test('remove from schema', () { + db.executeInTx( + 'SELECT powersync_replace_schema(?)', [json.encode(schema())]); db.execute( 'INSERT INTO users (id, name) VALUES (?, ?)', ['id', 'name']); db.executeInTx('SELECT powersync_replace_schema(?)', [ diff --git a/dart/test/sync_test.dart b/dart/test/sync_test.dart index 4aaf199..253ad66 100644 --- a/dart/test/sync_test.dart +++ b/dart/test/sync_test.dart @@ -2181,8 +2181,8 @@ CREATE TRIGGER users_ref_delete }); group('direct tables', () { - test('smoke test', () { - final schema = { + Object schema({Map additionalOptions = const {}}) { + return { 'tables': [ { 'name': 'users', @@ -2190,13 +2190,16 @@ CREATE TRIGGER users_ref_delete {'name': 'name', 'type': 'text'} ], 'direct': true, + ...additionalOptions, } ] }; + } + test('smoke test', () { db.executeInTx( - 'SELECT powersync_replace_schema(?)', [json.encode(schema)]); - invokeControl('start', json.encode({'schema': schema})); + 'SELECT powersync_replace_schema(?)', [json.encode(schema())]); + invokeControl('start', json.encode({'schema': schema()})); // Insert pushCheckpoint(buckets: [bucketDescription('a')]); @@ -2233,6 +2236,28 @@ CREATE TRIGGER users_ref_delete expect(db.select('SELECT * FROM users'), isEmpty); }); + + test('local only', () { + final localOnlySchema = schema(additionalOptions: {'local_only': true}); + + db.executeInTx( + 'SELECT powersync_replace_schema(?)', [json.encode(localOnlySchema)]); + invokeControl('start', json.encode({'schema': localOnlySchema})); + + // Insert + pushCheckpoint(buckets: [bucketDescription('a')]); + pushSyncData( + 'a', + '1', + 'my_user', + 'PUT', + {'name': 'First user'}, + objectType: 'users', + ); + pushCheckpointComplete(); + + expect(db.select('SELECT * FROM ps_untyped'), hasLength(1)); + }); }); test('can close database while iteration is active', () { From ff241b30d32e298992d097c682bed9ebd229350c Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Wed, 9 Sep 2026 14:55:12 +0200 Subject: [PATCH 6/7] Use _rest column pattern --- crates/core/src/json_util.rs | 6 ++- crates/core/src/schema/common.rs | 36 ++++++++-------- crates/core/src/schema/inspection.rs | 33 ++++++++++++++- crates/core/src/schema/management.rs | 27 +++++------- crates/core/src/schema/raw_table.rs | 26 ++++-------- crates/core/src/schema/table_info.rs | 61 ++++++++++++++++------------ crates/core/src/sync/mod.rs | 1 + crates/core/src/sync/sync_local.rs | 2 +- crates/core/src/utils/sql_buffer.rs | 4 +- crates/core/src/views.rs | 27 +++++------- dart/test/schema_test.dart | 10 +++-- dart/test/sync_test.dart | 2 +- 12 files changed, 130 insertions(+), 105 deletions(-) diff --git a/crates/core/src/json_util.rs b/crates/core/src/json_util.rs index 2d3fa04..6585181 100644 --- a/crates/core/src/json_util.rs +++ b/crates/core/src/json_util.rs @@ -6,8 +6,8 @@ use core::ffi::c_int; use crate::constants::SUBTYPE_JSON; use crate::create_sqlite_text_fn; use crate::error::{PowerSyncError, Result}; -use powersync_sqlite_nostd as sqlite; use powersync_sqlite_nostd::bindings::{SQLITE_RESULT_SUBTYPE, SQLITE_SUBTYPE}; +use powersync_sqlite_nostd::{self as sqlite, ColumnType}; use powersync_sqlite_nostd::{Connection, Context, Value}; use sqlite::ResultCode; @@ -38,6 +38,10 @@ fn powersync_json_merge_impl( } let mut result = String::from("{"); for arg in args { + if arg.value_type() == ColumnType::Null { + continue; + } + let chunk = arg.text(); if chunk.is_empty() || !chunk.starts_with('{') || !chunk.ends_with('}') { return Err(PowerSyncError::argument_error("Expected json object")); diff --git a/crates/core/src/schema/common.rs b/crates/core/src/schema/common.rs index 3600b22..c0d93d5 100644 --- a/crates/core/src/schema/common.rs +++ b/crates/core/src/schema/common.rs @@ -10,7 +10,7 @@ use serde::Deserialize; use crate::{ schema::{ Column, CommonTableOptions, PendingStatement, PendingStatementValue, RawTable, Table, - raw_table::InferredTableStructure, + raw_table::InferredTableStructure, table_info::RestColumnIndex, }, utils::SqlBuffer, }; @@ -37,14 +37,6 @@ impl<'a> SchemaTable<'a> { } } - pub fn data_column(&self) -> Option<&'static str> { - if let SchemaTable::Json(table) = self { - Some(table.data_column_name()) - } else { - None - } - } - pub fn common_options(&self) -> &CommonTableOptions { match self { Self::Json(table) => &table.options, @@ -75,13 +67,16 @@ impl<'a> SchemaTable<'a> { pub fn infer_put_stmt(&self, table_name: &str) -> PendingStatement { let mut buffer = SqlBuffer::new(); let mut params = vec![]; - let data_column = self.data_column(); + let mut rest = match self { + SchemaTable::Json(_) => Some(("_rest", RestColumnIndex::default())), + SchemaTable::Raw { .. } => None, + }; buffer.push_str("INSERT INTO "); let _ = buffer.identifier().write_str(table_name); buffer.push_str(" (id"); - if let Some(data_column) = data_column { - let _ = write!(&mut buffer, ", {data_column}"); + if let Some((column, _)) = rest { + let _ = write!(&mut buffer, ", {column}"); } for column in self.column_names() { @@ -90,23 +85,28 @@ impl<'a> SchemaTable<'a> { } buffer.push_str(") VALUES (?1"); params.push(PendingStatementValue::Id); - if data_column.is_some() { - params.push(PendingStatementValue::Row); + if let Some((_, ref mut rest_index)) = rest { + params.push(PendingStatementValue::Rest); buffer.push_str(", ?2"); + rest_index.rest_parameter_positions.push(1); // this is zero-indexed } - let data_start_index = if data_column.is_some() { 3 } else { 2 }; + let data_start_index = if rest.is_some() { 3 } else { 2 }; for (i, column) in self.column_names().enumerate() { buffer.comma(); let _ = write!(&mut buffer, "?{}", i + data_start_index); params.push(PendingStatementValue::Column(column.to_string())); + + if let Some((_, ref mut index)) = rest { + index.named_parameters.insert(column.to_string()); + } } buffer.push_str(") ON CONFLICT (id) DO UPDATE SET "); let mut do_update = buffer.comma_separated(); - if let Some(data_column) = data_column { + if let Some((column, _)) = rest { let entry = do_update.element(); - let _ = write!(entry, "{data_column} = ?2"); + let _ = write!(entry, "{column} = ?2"); } // Generate an "x" = ? for all synced columns to update them without affecting local-only @@ -120,7 +120,7 @@ impl<'a> SchemaTable<'a> { PendingStatement { sql: buffer.sql, params, - named_parameters_index: None, + named_parameters_index: rest.map(|e| e.1), } } diff --git a/crates/core/src/schema/inspection.rs b/crates/core/src/schema/inspection.rs index 902320d..181c60d 100644 --- a/crates/core/src/schema/inspection.rs +++ b/crates/core/src/schema/inspection.rs @@ -1,3 +1,5 @@ +use core::fmt::Write; + use alloc::borrow::ToOwned; use alloc::{format, vec}; use alloc::{string::String, vec::Vec}; @@ -6,6 +8,7 @@ use crate::error::Result; use crate::schema::raw_table::InferredTableStructure; use crate::utils::SqlBuffer; use crate::utils::database::Database; +use crate::views::table_columns_to_json_object; /// An existing PowerSync-managed view that was found in the schema. #[derive(PartialEq)] @@ -113,15 +116,16 @@ impl ExistingTable { local_only: local_only, direct: None, }); - } else if sql.contains("/* ps-managed */") && !ignore_direct { + } else if sql.contains("/* ps-managed") && !ignore_direct { results.push(ExistingTable { internal_name: internal_name.to_owned(), name: internal_name.to_owned(), - local_only: false, + local_only: sql.contains("local-only"), direct: Some(InferredTableStructure::read_from_database( internal_name, db, &None, + true, )?), }); } @@ -145,4 +149,29 @@ impl ExistingTable { None } } + + pub fn move_into_ps_untyped(&self, db: Database) -> Result<()> { + if self.local_only { + return Ok(()); + } + + let mut buffer = SqlBuffer::new(); + buffer.push_str("INSERT INTO ps_untyped(type, id, data) SELECT ?, id, "); + + if let Some(ref schema) = self.direct { + buffer.push_str("powersync_json_merge("); + buffer.push_str(&table_columns_to_json_object( + &self.internal_name, + &schema.columns, + )?); + buffer.push_str(", _rest)"); + } else { + buffer.push_str("data"); + } + + buffer.push_str(" FROM "); + let _ = buffer.identifier().write_str(&self.internal_name); + + db.exec_text(&buffer.sql, &self.name) + } } diff --git a/crates/core/src/schema/management.rs b/crates/core/src/schema/management.rs index 24a841d..03e3eb5 100644 --- a/crates/core/src/schema/management.rs +++ b/crates/core/src/schema/management.rs @@ -16,7 +16,7 @@ use sqlite::{Connection, ResultCode, Value}; use crate::create_sqlite_text_fn; use crate::error::{PowerSyncError, Result}; use crate::schema::inspection::{ExistingTable, ExistingView}; -use crate::schema::table_info::{Index, data_column_name}; +use crate::schema::table_info::Index; use crate::state::DatabaseState; use crate::utils::database::Database; use crate::utils::{SqlBuffer, verify_in_transaction}; @@ -55,24 +55,26 @@ fn update_tables(db: Database, schema: &Schema) -> Result<()> { } // New table. - let data_column = table.data_column_name(); let mut create_table = SqlBuffer::default(); create_table.push_str("CREATE TABLE "); table.write_name(&mut create_table); - _ = write!( - &mut create_table, - "(id TEXT PRIMARY KEY NOT NULL, {data_column} TEXT" - ); + _ = write!(&mut create_table, "(id TEXT PRIMARY KEY NOT NULL"); if table.direct { - create_table.push_str("/* ps-managed */"); + create_table.push_str(", _rest TEXT /* ps-managed "); + if table.local_only() { + create_table.push_str("local-only "); + } + create_table.push_str("*/"); for column in &table.columns { create_table.push_char(','); let _ = create_table.identifier().write_str(&column.name); let _ = write!(&mut create_table, " {}", column.type_name); } + } else { + create_table.push_str(", data TEXT"); } create_table.push_str(");"); db.exec_safe_str(&create_table.sql)?; @@ -86,16 +88,7 @@ fn update_tables(db: Database, schema: &Schema) -> Result<()> { // Remaining tables need to be dropped. But first, we want to move their contents to // ps_untyped. for remaining in existing_tables.values() { - if !remaining.local_only { - db.exec_text( - &format!( - "INSERT INTO ps_untyped(type, id, data) SELECT ?, id, {} FROM {:}", - data_column_name(remaining.direct.is_some()), - SqlBuffer::quote_identifier(&remaining.internal_name) - ), - &remaining.name, - )?; - } + remaining.move_into_ps_untyped(db)?; } // We cannot have any open queries on sqlite_master at the point that we drop tables, otherwise diff --git a/crates/core/src/schema/raw_table.rs b/crates/core/src/schema/raw_table.rs index 54dbdec..57631b7 100644 --- a/crates/core/src/schema/raw_table.rs +++ b/crates/core/src/schema/raw_table.rs @@ -25,6 +25,7 @@ impl InferredTableStructure { table_name: &str, db: Database, synced_columns: &Option, + is_direct: bool, ) -> Result { let stmt = db.prepare_v2("select name, type from pragma_table_info(?)")?; stmt.bind_text(1, table_name, Destructor::STATIC)?; @@ -42,6 +43,8 @@ impl InferredTableStructure { && !filter.matches(name) { // This column isn't part of the synced columns, skip. + } else if is_direct && name == "_rest" { + // _rest column is an artifact of direct tables, skip. } else { columns.push(Column { name: name.to_owned(), @@ -136,6 +139,7 @@ impl SchemaCacheEntry { local_table_name, db, &table.schema.synced_columns, + false, )?; let schema_table = SchemaTable::Raw { definition: table, @@ -161,7 +165,7 @@ pub fn generate_raw_table_trigger( let local_table_name = table.require_table_name()?; let synced_columns = &table.schema.synced_columns; let resolved_table = - InferredTableStructure::read_from_database(local_table_name, db, synced_columns)?; + InferredTableStructure::read_from_database(local_table_name, db, synced_columns, false)?; let as_schema_table = SchemaTable::Raw { definition: table, @@ -218,18 +222,19 @@ pub fn generate_schema_table_trigger( buffer.push_str("SELECT RAISE(FAIL, 'Unexpected update on insert-only table');\n"); } else if !flags.local_only() { // Insert-only tables use manual CRUD writes so they don't block incoming data. - let fragment = table_columns_to_json_object("NEW", &table)?; + let fragment = table_columns_to_json_object("NEW", table.columns())?; buffer.powersync_crud_manual_put(table.name(), &fragment); } } else { if write == WriteType::Update { // Updates must not change the id. buffer.check_id_not_changed(); + has_stmt = true; } - let json_fragment_new = table_columns_to_json_object("NEW", &table)?; + let json_fragment_new = table_columns_to_json_object("NEW", table.columns())?; let json_fragment_old = if write == WriteType::Update { - Some(table_columns_to_json_object("OLD", &table)?) + Some(table_columns_to_json_object("OLD", table.columns())?) } else { None }; @@ -248,19 +253,6 @@ pub fn generate_schema_table_trigger( write!(f, ", {json_fragment_new}))") }); - if write != WriteType::Delete - && let Some(data_column) = table.data_column() - { - // If the table has a __data column storing the full JSON row, we also need to update - // that. - let _ = write!( - &mut buffer, - "UPDATE {local_table_name} SET {data_column} = {json_fragment_new} WHERE id = NEW.id;\n" - ); - - has_stmt = true; - } - if !flags.local_only() { has_stmt = true; buffer.insert_into_powersync_crud(InsertIntoCrud { diff --git a/crates/core/src/schema/table_info.rs b/crates/core/src/schema/table_info.rs index f4d3f0e..8de9984 100644 --- a/crates/core/src/schema/table_info.rs +++ b/crates/core/src/schema/table_info.rs @@ -4,11 +4,13 @@ use alloc::rc::Rc; use alloc::string::ToString; use alloc::vec; use alloc::{collections::btree_set::BTreeSet, format, string::String, vec::Vec}; +use powersync_sqlite_nostd::Destructor; use serde::{Deserialize, de::Visitor}; use crate::error::PowerSyncError; use crate::schema::raw_table::generate_schema_table_trigger; use crate::schema::{ColumnFilter, SchemaTable}; +use crate::sync::PreparedPendingStatement; use crate::utils::database::Database; use crate::utils::{SqlBuffer, WriteType}; @@ -87,32 +89,44 @@ impl Table { } pub fn move_from_ps_untyped(&self, db: Database) -> Result<(), PowerSyncError> { - let mut stmt = SqlBuffer::default(); let direct = self.direct; - stmt.push_str("INSERT INTO "); - self.write_name(&mut stmt); - let _ = write!(&mut stmt, "(id, {}", self.data_column_name()); + let mut delete_stmt = SqlBuffer::new(); + delete_stmt.push_str("DELETE FROM ps_untyped WHERE type = ?"); if direct { - for column in &self.columns { - stmt.push_char(','); - let _ = stmt.identifier().write_str(&column.name); - } - } - - stmt.push_str(") SELECT id, data"); - if direct { - for column in &self.columns { - stmt.push_char(','); - stmt.json_extract_and_cast("data", &column.name, &column.type_name); + // Copying into direct tables reqires extracting from JSON. This essentially replays a + // sync_local step for the table, using ps_untyped as source. + let stmt = Rc::new(SchemaTable::Json(self).infer_put_stmt(&self.name)); + let stmt = PreparedPendingStatement::prepare(db, stmt)?; + + let _ = delete_stmt.write_str(" RETURNING id, data"); + let source = db.prepare_v2(&delete_stmt.sql)?; + source.bind_text(1, &self.name, Destructor::STATIC)?; + + while source.step()? { + let id = source.column_text(0)?; + let data = source.column_text(1)?; + + let parsed: serde_json::Value = + serde_json::from_str(data).map_err(PowerSyncError::json_local_error)?; + let json_object = parsed.as_object().ok_or_else(|| { + PowerSyncError::argument_error("expected oplog data to be an object") + })?; + let rest = stmt.render_rest_object(json_object)?; + stmt.bind_for_put(id, data, Some(json_object), rest.as_ref())?; + stmt.exec(&self.name, id, Some(&data))?; } + } else { + let mut stmt = SqlBuffer::default(); + stmt.push_str("INSERT INTO "); + self.write_name(&mut stmt); + let _ = stmt.write_str(" (id, data) SELECT id, data FROM ps_untyped WHERE type = ?"); + let _ = db.exec_text(&stmt.sql, &self.name); + db.exec_text(&delete_stmt.sql, &self.name)?; } - stmt.push_str(" FROM ps_untyped WHERE type = ?"); - - db.exec_text(&stmt.sql, &self.name)?; - db.exec_text("DELETE FROM ps_untyped WHERE type = ?", &self.name) + Ok(()) } pub fn write_name(&self, buffer: &mut SqlBuffer) { @@ -124,10 +138,6 @@ impl Table { } } - pub fn data_column_name(&self) -> &'static str { - data_column_name(self.direct) - } - pub fn generate_direct_trigger(&self, write: WriteType) -> Result { debug_assert!(self.direct); generate_schema_table_trigger( @@ -140,10 +150,6 @@ impl Table { } } -pub fn data_column_name(is_direct: bool) -> &'static str { - if is_direct { "__data" } else { "data" } -} - impl RawTable { pub fn require_table_name(&self) -> Result<&str, PowerSyncError> { let Some(local_table_name) = self.schema.table_name.as_ref() else { @@ -367,6 +373,7 @@ pub struct PendingStatement { pub named_parameters_index: Option, } +#[derive(Default)] pub struct RestColumnIndex { /// All column names referenced by this statement. pub named_parameters: BTreeSet, diff --git a/crates/core/src/sync/mod.rs b/crates/core/src/sync/mod.rs index 874eb29..2c150de 100644 --- a/crates/core/src/sync/mod.rs +++ b/crates/core/src/sync/mod.rs @@ -19,6 +19,7 @@ pub use checksum::Checksum; use crate::state::DatabaseState; pub use streaming_sync::SyncClient; +pub use sync_local::PreparedPendingStatement; pub fn register(db: *mut sqlite::sqlite3, state: Rc) -> Result<(), ResultCode> { interface::register(db, state) diff --git a/crates/core/src/sync/sync_local.rs b/crates/core/src/sync/sync_local.rs index 24f6b19..01a3ac7 100644 --- a/crates/core/src/sync/sync_local.rs +++ b/crates/core/src/sync/sync_local.rs @@ -476,7 +476,7 @@ enum TableDefinition<'a> { Direct(&'a Table), } -struct PreparedPendingStatement { +pub struct PreparedPendingStatement { stmt: Statement, needs_parsed_json: bool, definition: Rc, diff --git a/crates/core/src/utils/sql_buffer.rs b/crates/core/src/utils/sql_buffer.rs index 6a9c97d..3c68e24 100644 --- a/crates/core/src/utils/sql_buffer.rs +++ b/crates/core/src/utils/sql_buffer.rs @@ -123,7 +123,7 @@ impl SqlBuffer { Some(include_old) => { let old_values = table_columns_to_json_object_with_filter( "OLD", - insert.table, + insert.table.columns(), include_old.column_filter(), )?; @@ -134,7 +134,7 @@ impl SqlBuffer { // only include the powersync_diff of columns matched by the filter. let filtered_new_fragment = table_columns_to_json_object_with_filter( "NEW", - insert.table, + insert.table.columns(), include_old.column_filter(), )?; diff --git a/crates/core/src/views.rs b/crates/core/src/views.rs index f5e7089..ffe2e3e 100644 --- a/crates/core/src/views.rs +++ b/crates/core/src/views.rs @@ -6,7 +6,7 @@ use core::fmt::{Write, from_fn}; use core::mem; use crate::error::{PowerSyncError, Result}; -use crate::schema::{ColumnFilter, SchemaTable, Table}; +use crate::schema::{Column, ColumnFilter, SchemaTable, Table}; use crate::utils::{InsertIntoCrud, SqlBuffer, WriteType}; pub fn powersync_view_sql(table_info: &Table) -> String { @@ -140,7 +140,7 @@ pub fn powersync_trigger_insert_sql(table_info: &Table) -> Result { sql.check_id_valid(); } - let json_fragment = table_columns_to_json_object("NEW", &as_schema_table)?; + let json_fragment = table_columns_to_json_object("NEW", &table_info.columns)?; if insert_only { // This is using the manual powersync_crud_ instead of powersync_crud because insert-only @@ -188,7 +188,6 @@ pub fn powersync_trigger_update_sql(table_info: &Table) -> Result { let name = &table_info.name; let view_name = table_info.view_name(); let local_only = table_info.options.flags.local_only(); - let as_schema_table = SchemaTable::from(table_info); let mut sql = SqlBuffer::new(); sql.create_trigger("ps_view_update_", view_name); @@ -202,8 +201,8 @@ pub fn powersync_trigger_update_sql(table_info: &Table) -> Result { sql.push_str("BEGIN\n"); sql.check_id_not_changed(); - let json_fragment_new = table_columns_to_json_object("NEW", &as_schema_table)?; - let json_fragment_old = table_columns_to_json_object("OLD", &as_schema_table)?; + let json_fragment_new = table_columns_to_json_object("NEW", &table_info.columns)?; + let json_fragment_old = table_columns_to_json_object("OLD", &table_info.columns)?; // UPDATE {internal_name} SET data = {json_fragment_new} WHERE id = NEW.id; sql.push_str("UPDATE "); @@ -218,7 +217,7 @@ pub fn powersync_trigger_update_sql(table_info: &Table) -> Result { sql.insert_into_powersync_crud(InsertIntoCrud { op: WriteType::Update, id_expr: "NEW.id", - table: &as_schema_table, + table: &SchemaTable::Json(table_info), type_name: name, data: Some(&from_fn(|f| { write!( @@ -241,16 +240,13 @@ pub fn powersync_trigger_update_sql(table_info: &Table) -> Result { /// Given a query returning column names, return a JSON object fragment for a trigger. /// /// Example output with prefix "NEW": "json_object('id', NEW.id, 'name', NEW.name, 'age', NEW.age)". -pub fn table_columns_to_json_object<'a>( - prefix: &str, - table: &'a SchemaTable<'a>, -) -> Result { - table_columns_to_json_object_with_filter(prefix, table, None) +pub fn table_columns_to_json_object(prefix: &str, columns: &[Column]) -> Result { + table_columns_to_json_object_with_filter(prefix, columns, None) } pub fn table_columns_to_json_object_with_filter<'a>( prefix: &str, - table: &'a SchemaTable<'a>, + columns: &[Column], filter: Option<&'a ColumnFilter>, ) -> Result { // floor(SQLITE_MAX_FUNCTION_ARG / 2). @@ -274,8 +270,7 @@ pub fn table_columns_to_json_object_with_filter<'a>( buffer.sql } - let mut columns = table.column_names(); - while let Some(name) = columns.next() { + for Column { name, type_name: _ } in columns { if let Some(filter) = filter && !filter.matches(name) { @@ -377,8 +372,8 @@ mod test { #[test] fn test_json_object_fragment() { - let fragment = - table_columns_to_json_object("NEW", &(&test_table()).into()).expect("should generate"); + let columns = &test_table().columns; + let fragment = table_columns_to_json_object("NEW", columns).expect("should generate"); assert_eq!( fragment, diff --git a/dart/test/schema_test.dart b/dart/test/schema_test.dart index f51106f..c14659e 100644 --- a/dart/test/schema_test.dart +++ b/dart/test/schema_test.dart @@ -355,7 +355,7 @@ END''', { 'id': 'user-id', 'name': 'Name', - '__data': '{"name":"Name","other":3}' + '_rest': '{"other":3}', }, ]); @@ -365,7 +365,7 @@ END''', )[0].columnAt(0); expect( createTable, - 'CREATE TABLE "users"(id TEXT PRIMARY KEY NOT NULL, __data TEXT/* ps-managed */,"name" text)', + 'CREATE TABLE "users"(id TEXT PRIMARY KEY NOT NULL, _rest TEXT /* ps-managed */,"name" text)', ); final triggers = db @@ -388,7 +388,6 @@ END''', r''' CREATE TRIGGER "users_trigger_UPDATE" AFTER UPDATE ON "users" FOR EACH ROW WHEN NOT powersync_in_sync_operation() BEGIN SELECT CASE WHEN (OLD.id != NEW.id) THEN RAISE (FAIL, 'Cannot update id') END; -UPDATE users SET __data = json_object('name', powersync_strip_subtype(NEW."name")) WHERE id = NEW.id; INSERT INTO powersync_crud(op,id,type,data,options) VALUES ('PATCH', NEW.id, 'users', json(powersync_diff(json_object('name', powersync_strip_subtype(OLD."name")), json_object('name', powersync_strip_subtype(NEW."name")))), 0); END''' ]); @@ -416,6 +415,11 @@ END''' expect(db.select('SELECT * FROM ps_untyped'), [ {'type': 'users', 'id': 'id', 'data': '{"name":"name"}'} ]); + + expect( + db.select( + 'SELECT * FROM sqlite_schema WHERE type = ?', ['trigger']), + isEmpty); }); }); }); diff --git a/dart/test/sync_test.dart b/dart/test/sync_test.dart index 253ad66..211a205 100644 --- a/dart/test/sync_test.dart +++ b/dart/test/sync_test.dart @@ -2218,7 +2218,7 @@ CREATE TRIGGER users_ref_delete { 'id': 'my_user', 'name': 'First user', - '__data': '{"name":"First user"}' + '_rest': null, } ]); From 9542a89d5d48f1ecc028d577f342ff56920b76a8 Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Thu, 10 Sep 2026 16:48:17 +0200 Subject: [PATCH 7/7] Support migrating from JSON tables --- crates/core/src/migrations.rs | 9 ++- crates/core/src/schema/inspection.rs | 6 +- crates/core/src/schema/management.rs | 105 +++++++++++++++++++-------- crates/core/src/schema/table_info.rs | 62 +++++++++++----- crates/core/src/view_admin.rs | 7 +- dart/test/schema_test.dart | 72 +++++++++++++++--- 6 files changed, 195 insertions(+), 66 deletions(-) diff --git a/crates/core/src/migrations.rs b/crates/core/src/migrations.rs index 3a8a58e..b1c72ff 100644 --- a/crates/core/src/migrations.rs +++ b/crates/core/src/migrations.rs @@ -4,7 +4,6 @@ use alloc::format; use alloc::string::{String, ToString}; use alloc::vec::Vec; -use powersync_sqlite_nostd::Context; use powersync_sqlite_nostd::{self as sqlite, Destructor}; use serde::Serialize; use serde_json::json; @@ -15,12 +14,16 @@ use crate::fix_data::apply_v035_fix; use crate::schema::inspection::ExistingView; use crate::sync::BucketPriority; use crate::utils::database::Database; +use crate::utils::verify_in_transaction; pub const LATEST_VERSION: i32 = 14; -pub fn powersync_migrate(ctx: *mut sqlite::context, target_version: i32) -> Result<()> { - let local_db = Database::from(ctx.db_handle()); +pub fn initialize_database(db: Database) -> Result<()> { + verify_in_transaction(db)?; + powersync_migrate(db, LATEST_VERSION) +} +pub fn powersync_migrate(local_db: Database, target_version: i32) -> Result<()> { // language=SQLite local_db.exec_safe( c"\ diff --git a/crates/core/src/schema/inspection.rs b/crates/core/src/schema/inspection.rs index 181c60d..3e76857 100644 --- a/crates/core/src/schema/inspection.rs +++ b/crates/core/src/schema/inspection.rs @@ -17,7 +17,7 @@ pub struct ExistingView { pub name: String, /// SQL contents of the `CREATE VIEW` statement. /// - /// This is not set for as_raw_table tables, which don't have a view. + /// This is not set for direct tables, which don't have a view. pub sql: Option, /// SQL contents of all triggers implementing deletes by forwarding to /// `ps_data` and `ps_crud`. @@ -74,6 +74,10 @@ SELECT Ok(()) } + pub fn delete_from_db(&self, db: Database) -> Result<()> { + Self::drop_by_name(db, &self.name) + } + pub fn create(&self, db: Database) -> Result<()> { if let Some(create_view) = &self.sql { Self::drop_by_name(db, &self.name)?; diff --git a/crates/core/src/schema/management.rs b/crates/core/src/schema/management.rs index 03e3eb5..75aa052 100644 --- a/crates/core/src/schema/management.rs +++ b/crates/core/src/schema/management.rs @@ -15,6 +15,7 @@ use sqlite::{Connection, ResultCode, Value}; use crate::create_sqlite_text_fn; use crate::error::{PowerSyncError, Result}; +use crate::migrations::initialize_database; use crate::schema::inspection::{ExistingTable, ExistingView}; use crate::schema::table_info::Index; use crate::state::DatabaseState; @@ -27,7 +28,11 @@ use crate::views::{ use super::Schema; -fn update_tables(db: Database, schema: &Schema) -> Result<()> { +fn update_tables( + db: Database, + schema: &Schema, + existing_views: &mut BTreeMap<&str, &ExistingView>, +) -> Result<()> { let existing_tables = ExistingTable::list(db)?; let mut existing_tables = { let mut map = BTreeMap::new(); @@ -38,19 +43,57 @@ fn update_tables(db: Database, schema: &Schema) -> Result<()> { }; for table in &schema.tables { + let mut move_data_from = None::<&str>; + if let Some(existing) = existing_tables.remove(&*table.name) { - if !table.direct && existing.local_only != table.local_only() { - // Migrating between local-only and synced tables. This works by deleting - // existing and re-creating the table from scratch. We can re-create first and - // delete the old table afterwards because they have a different name - // (local-only tables have a ps_data_local prefix). - // Direct tables are the same whether they're local or not. - - // To delete the old existing table in the end. - existing_tables.insert(&existing.name, existing); - } else { - // Compatible table exists already, nothing to do. - continue; + match (&existing.direct, table.direct) { + (None, false) => { + // JSON-based table before and now. We might have to migrate between synced and + // local-only tables. + if existing.local_only != table.local_only() { + // Migrating between local-only and synced tables. This works by deleting + // existing and re-creating the table from scratch. We can re-create first + // and delete the old table afterwards because they have a different name + // (local-only tables have a ps_data_local prefix). + + // To delete the old existing table in the end. + existing_tables.insert(&existing.name, existing); + } else { + // Compatible table exists already, nothing to do. + continue; + } + } + (None, true) => { + // When migrating from JSON-based to direct tables, there are four cases to + // consider: + // 1. Local-only to direct local-only: We copy data; delete the old table. + // 2. Local-only to synced: Delete old table, copy from ps_untyped for new. + // 3. Synced to local-only: Move old into ps_untyped; create new from scratch. + // 4. Synced to synced: Copy data; delete old table. + if existing.local_only == table.local_only() { + // Case 1 or 4. + move_data_from = Some(&existing.internal_name); + } else { + // Case 2 and 3 is the default, we'll delete the old table in the end which + // moves to ps_untyped if necessary. + } + + // To delete the existing table in the end. + existing_tables.insert(&existing.name, existing); + + // The direct table we create conflicts with the view. So delete that one first. + if let Some(old_view) = existing_views.remove(&*existing.name) { + old_view.delete_from_db(db)?; + } + } + (Some(_), false) => { + return Err(PowerSyncError::argument_error( + "Switching from direct to json-based tables is not yet implemented.", + )); + } + (Some(_), true) => { + // TODO: Consider migrations in schema tables. + } } } @@ -79,7 +122,9 @@ fn update_tables(db: Database, schema: &Schema) -> Result<()> { create_table.push_str(");"); db.exec_safe_str(&create_table.sql)?; - if !table.local_only() { + if let Some(old_json_table) = move_data_from { + table.direct_move_from_json(db, old_json_table)?; + } else if !table.local_only() { // MOVE data if any table.move_from_ps_untyped(db)?; } @@ -203,17 +248,11 @@ SELECT Ok(()) } -fn update_views(db: Database, schema: &Schema) -> Result<()> { - // First, find all existing views and index them by name. - let existing = ExistingView::list(db)?; - let mut existing = { - let mut map = BTreeMap::new(); - for entry in &existing { - map.insert(&*entry.name, entry); - } - map - }; - +fn update_views( + db: Database, + schema: &Schema, + existing: &mut BTreeMap<&str, &ExistingView>, +) -> Result<()> { for table in &schema.tables { let view_sql = if table.direct { None @@ -267,12 +306,20 @@ fn powersync_replace_schema_impl( let parsed_schema = serde_json::from_str::(schema).map_err(PowerSyncError::as_argument_error)?; - // language=SQLite - db.exec_safe(c"SELECT powersync_init()")?; + initialize_database(db)?; + + let views = ExistingView::list(db)?; + let mut existing_views = { + let mut map = BTreeMap::new(); + for entry in &views { + map.insert(&*entry.name, entry); + } + map + }; - update_tables(db, &parsed_schema)?; + update_tables(db, &parsed_schema, &mut existing_views)?; update_indexes(db, &parsed_schema)?; - update_views(db, &parsed_schema)?; + update_views(db, &parsed_schema, &mut existing_views)?; state.set_schema(parsed_schema); Ok(String::from("")) diff --git a/crates/core/src/schema/table_info.rs b/crates/core/src/schema/table_info.rs index 8de9984..f745958 100644 --- a/crates/core/src/schema/table_info.rs +++ b/crates/core/src/schema/table_info.rs @@ -11,7 +11,7 @@ use crate::error::PowerSyncError; use crate::schema::raw_table::generate_schema_table_trigger; use crate::schema::{ColumnFilter, SchemaTable}; use crate::sync::PreparedPendingStatement; -use crate::utils::database::Database; +use crate::utils::database::{Database, Statement}; use crate::utils::{SqlBuffer, WriteType}; #[derive(Deserialize)] @@ -95,28 +95,11 @@ impl Table { delete_stmt.push_str("DELETE FROM ps_untyped WHERE type = ?"); if direct { - // Copying into direct tables reqires extracting from JSON. This essentially replays a - // sync_local step for the table, using ps_untyped as source. - let stmt = Rc::new(SchemaTable::Json(self).infer_put_stmt(&self.name)); - let stmt = PreparedPendingStatement::prepare(db, stmt)?; - let _ = delete_stmt.write_str(" RETURNING id, data"); let source = db.prepare_v2(&delete_stmt.sql)?; source.bind_text(1, &self.name, Destructor::STATIC)?; - while source.step()? { - let id = source.column_text(0)?; - let data = source.column_text(1)?; - - let parsed: serde_json::Value = - serde_json::from_str(data).map_err(PowerSyncError::json_local_error)?; - let json_object = parsed.as_object().ok_or_else(|| { - PowerSyncError::argument_error("expected oplog data to be an object") - })?; - let rest = stmt.render_rest_object(json_object)?; - stmt.bind_for_put(id, data, Some(json_object), rest.as_ref())?; - stmt.exec(&self.name, id, Some(&data))?; - } + self.direct_move_from_stmt(db, source)?; } else { let mut stmt = SqlBuffer::default(); stmt.push_str("INSERT INTO "); @@ -129,6 +112,47 @@ impl Table { Ok(()) } + pub fn direct_move_from_json( + &self, + db: Database, + json_table: &str, + ) -> Result<(), PowerSyncError> { + debug_assert!(self.direct); + + let mut source = SqlBuffer::new(); + source.push_str("SELECT id, data FROM "); + let _ = write!(source.identifier(), "{}", json_table); + + let source = db.prepare_v2(&source.sql)?; + self.direct_move_from_stmt(db, source) + } + + /// For direct tables, copies data from a prepared statement returning id and data. + fn direct_move_from_stmt(&self, db: Database, source: Statement) -> Result<(), PowerSyncError> { + debug_assert!(self.direct); + + // Copying into direct tables reqires extracting from JSON. This essentially replays a + // sync_local step for the table, using a custom source. + let stmt = Rc::new(SchemaTable::Json(self).infer_put_stmt(&self.name)); + let stmt = PreparedPendingStatement::prepare(db, stmt)?; + + while source.step()? { + let id = source.column_text(0)?; + let data = source.column_text(1)?; + + let parsed: serde_json::Value = + serde_json::from_str(data).map_err(PowerSyncError::json_local_error)?; + let json_object = parsed.as_object().ok_or_else(|| { + PowerSyncError::argument_error("expected oplog data to be an object") + })?; + let rest = stmt.render_rest_object(json_object)?; + stmt.bind_for_put(id, data, Some(json_object), rest.as_ref())?; + stmt.exec(&self.name, id, Some(&data))?; + } + + Ok(()) + } + pub fn write_name(&self, buffer: &mut SqlBuffer) { if self.direct { // Direct tables don't have views, so use the name of the table directly. diff --git a/crates/core/src/view_admin.rs b/crates/core/src/view_admin.rs index cfe8b93..702fa16 100644 --- a/crates/core/src/view_admin.rs +++ b/crates/core/src/view_admin.rs @@ -12,7 +12,7 @@ use sqlite::{ResultCode, Value}; use crate::create_sqlite_text_fn; use crate::error::{PowerSyncError, Result}; -use crate::migrations::{LATEST_VERSION, powersync_migrate}; +use crate::migrations::{initialize_database, powersync_migrate}; use crate::schema::inspection::ExistingView; use crate::state::DatabaseState; use crate::utils::database::Database; @@ -34,8 +34,7 @@ extern "C" fn powersync_drop_view( fn powersync_init_impl(ctx: *mut sqlite::context, _args: &[*mut sqlite::value]) -> Result { let db = Database::from(ctx.db_handle()); - verify_in_transaction(db)?; - powersync_migrate(ctx, LATEST_VERSION)?; + initialize_database(db)?; Ok(String::from("")) } @@ -50,7 +49,7 @@ fn powersync_test_migration_impl( verify_in_transaction(db)?; let target_version = args[0].int(); - powersync_migrate(ctx, target_version)?; + powersync_migrate(db, target_version)?; Ok(String::from("")) } diff --git a/dart/test/schema_test.dart b/dart/test/schema_test.dart index c14659e..fd395f5 100644 --- a/dart/test/schema_test.dart +++ b/dart/test/schema_test.dart @@ -339,17 +339,19 @@ END''', }; } + void replaceSchema(Object schema) { + db.executeInTx( + 'SELECT powersync_replace_schema(?)', [json.encode(schema)]); + } + test('create', () { - db.executeInTx('SELECT powersync_replace_schema(?)', [ - json.encode({'tables': []}) - ]); + replaceSchema({'tables': []}); db.execute('INSERT INTO ps_untyped (type, id, data) VALUES (?, ?, ?)', [ 'users', 'user-id', json.encode({'name': 'Name', 'other': 3}) ]); - db.executeInTx( - 'SELECT powersync_replace_schema(?)', [json.encode(schema())]); + replaceSchema(schema()); expect(db.select('SELECT * FROM users'), [ { @@ -394,9 +396,7 @@ END''' }); test('local-only', () { - db.executeInTx('SELECT powersync_replace_schema(?)', [ - json.encode(schema(additionalOptions: {'local_only': true})) - ]); + replaceSchema(schema(additionalOptions: {'local_only': true})); db.execute( 'INSERT INTO users (id, name) VALUES (?, ?)', ['id', 'name']); @@ -404,8 +404,7 @@ END''' }); test('remove from schema', () { - db.executeInTx( - 'SELECT powersync_replace_schema(?)', [json.encode(schema())]); + replaceSchema(schema()); db.execute( 'INSERT INTO users (id, name) VALUES (?, ?)', ['id', 'name']); db.executeInTx('SELECT powersync_replace_schema(?)', [ @@ -421,6 +420,59 @@ END''' 'SELECT * FROM sqlite_schema WHERE type = ?', ['trigger']), isEmpty); }); + + group('migrate', () { + group('from json to direct', () { + test('local-only', () { + replaceSchema(schema( + additionalOptions: {'local_only': true, 'direct': false})); + db.execute( + 'INSERT INTO users (id, name) VALUES (?, ?)', ['id', 'name']); + replaceSchema(schema(additionalOptions: {'local_only': true})); + expect(db.select('SELECT * FROM users'), hasLength(1)); + }); + + test('local-only to synced', () { + replaceSchema(schema( + additionalOptions: {'local_only': true, 'direct': false})); + db.execute( + 'INSERT INTO users (id, name) VALUES (?, ?)', ['id', 'name']); + replaceSchema(schema(additionalOptions: {})); + + // Migrating from local-only to synced tables deletes data + expect(db.select('SELECT * FROM users'), isEmpty); + }); + + test('synced', () { + replaceSchema(schema(additionalOptions: {'direct': false})); + db.execute( + 'INSERT INTO users (id, name) VALUES (?, ?)', ['id', 'name']); + replaceSchema(schema(additionalOptions: {})); + expect(db.select('SELECT * FROM users'), hasLength(1)); + expect(db.select('SELECT * FROM ps_crud'), hasLength(1)); + }); + + test('synced to local-only', () { + replaceSchema(schema(additionalOptions: {'direct': false})); + db.execute( + 'INSERT INTO users (id, name) VALUES (?, ?)', ['id', 'name']); + + replaceSchema(schema(additionalOptions: {'local_only': true})); + // Data should be deleted when changing to a local-only table, + // previous crud entry is still there. + expect(db.select('SELECT * FROM users'), isEmpty); + expect(db.select('SELECT * FROM ps_crud'), hasLength(1)); + }); + }); + + // todo: from json to direct + // todo: from direct to json + + // todo: add column + // todo: change column type + // todo: remove column + // todo: split columns + }); }); }); }