From d66189032075b29ca93c08914e93e4f197d03c66 Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Wed, 2 Sep 2026 13:57:14 +0200 Subject: [PATCH 1/3] Avoid passing partial arguments as string to sync_local --- crates/core/src/sync/storage_adapter.rs | 56 ++++++------------------- crates/core/src/sync/streaming_sync.rs | 13 ++++++ crates/core/src/sync/sync_local.rs | 38 ++++++++++------- 3 files changed, 50 insertions(+), 57 deletions(-) diff --git a/crates/core/src/sync/storage_adapter.rs b/crates/core/src/sync/storage_adapter.rs index c79ccf9..30084fe 100644 --- a/crates/core/src/sync/storage_adapter.rs +++ b/crates/core/src/sync/storage_adapter.rs @@ -2,7 +2,6 @@ use core::fmt::Display; use alloc::{rc::Rc, string::ToString, vec::Vec}; use powersync_sqlite_nostd::{self as sqlite}; -use serde::Serialize; use crate::{ error::{PowerSyncError, Result}, @@ -273,51 +272,22 @@ WHERE bucket = ?1", self.persist_last_seen_checkpoint_request_id(*checkpoint_request_id)?; } - #[derive(Serialize)] - struct PartialArgs<'a> { - priority: BucketPriority, - buckets: Vec<&'a str>, - } let now = self.now()?; - let sync_result = match priority { - None => { - let mut sync = SyncOperation::new(state, self.db, None, now); - sync.use_schema(schema); - sync.apply() - } - Some(priority) => { - let args = PartialArgs { + let mut sync = match priority { + None => SyncOperation::new(state, self.db, None, now), + Some(priority) => SyncOperation::new( + state, + self.db, + Some(PartialSyncOperation { priority, - buckets: checkpoint - .buckets - .values() - .filter_map(|item| { - if item.is_in_priority(Some(priority)) { - Some(item.bucket.as_str()) - } else { - None - } - }) - .collect(), - }; - - // TODO: Avoid this serialization, it's currently used to bind JSON SQL parameters. - let serialized_args = - serde_json::to_string(&args).map_err(PowerSyncError::internal)?; - let mut sync = SyncOperation::new( - state, - self.db, - Some(PartialSyncOperation { - priority, - args: &serialized_args, - }), - now, - ); - sync.use_schema(schema); - sync.apply() - } - }?; + checkpoint, + }), + now, + ), + }; + sync.use_schema(schema); + let sync_result = sync.apply()?; if sync_result == 1 { if priority.is_none() { diff --git a/crates/core/src/sync/streaming_sync.rs b/crates/core/src/sync/streaming_sync.rs index 2cddd9a..6994bb4 100644 --- a/crates/core/src/sync/streaming_sync.rs +++ b/crates/core/src/sync/streaming_sync.rs @@ -1009,6 +1009,19 @@ impl OwnedCheckpoint { self.last_op_id = diff.last_op_id; self.write_checkpoint = diff.write_checkpoint; } + + pub fn list_buckets<'a>( + &'a self, + min_priority: Option, + ) -> impl Iterator { + self.buckets.values().filter_map(move |item| { + if item.is_in_priority(min_priority) { + Some(item.bucket.as_str()) + } else { + None + } + }) + } } /// A transition representing pending changes between [StreamingSyncIteration::prepare_handling_sync_line] diff --git a/crates/core/src/sync/sync_local.rs b/crates/core/src/sync/sync_local.rs index c973412..19b46f0 100644 --- a/crates/core/src/sync/sync_local.rs +++ b/crates/core/src/sync/sync_local.rs @@ -2,6 +2,7 @@ use alloc::collections::btree_map::BTreeMap; use alloc::format; use alloc::rc::Rc; use alloc::string::{String, ToString}; +use alloc::vec::Vec; use serde::Serialize; use serde::ser::SerializeMap; @@ -15,6 +16,7 @@ use crate::sync::BucketPriority; use crate::sync::storage_adapter::{ LAST_SEEN_CHECKPOINT_REQUEST_ID_KEY, TARGET_CHECKPOINT_REQUEST_ID_KEY, }; +use crate::sync::streaming_sync::OwnedCheckpoint; use crate::sync::sync_status::TimestampMicros; use crate::utils::SqlBuffer; use crate::utils::database::{Database, Statement}; @@ -24,9 +26,13 @@ use powersync_sqlite_nostd::{self as sqlite, Destructor}; pub struct PartialSyncOperation<'a> { /// The lowest priority part of the partial sync operation. pub priority: BucketPriority, - /// The JSON-encoded arguments passed by the client SDK. This includes the priority and a list - /// of bucket names in that (and higher) priorities. - pub args: &'a str, + pub checkpoint: &'a OwnedCheckpoint, +} + +impl<'a> PartialSyncOperation<'a> { + fn list_buckets(&self) -> impl Iterator { + self.checkpoint.list_buckets(Some(self.priority)) + } } pub struct SyncOperation<'a> { @@ -285,8 +291,8 @@ SELECT -- We filter out duplicates using the GROUP BY below. WITH involved_buckets (id) AS MATERIALIZED ( - SELECT id FROM ps_buckets WHERE ?1 IS NULL - OR name IN (SELECT value FROM json_each(json_extract(?1, '$.buckets'))) + SELECT id FROM ps_buckets + WHERE name IN (SELECT value FROM json_each(?1)) ), updated_rows AS ( SELECT b.row_type, b.row_id FROM ps_buckets AS buckets @@ -314,7 +320,10 @@ SELECT -- Group for (2) GROUP BY b.row_type, b.row_id;", )?; - stmt.bind_text(1, partial.args, Destructor::STATIC)?; + + let bucket_ids: Vec<&str> = partial.list_buckets().collect(); + let bucket_ids = serde_json::to_string(&bucket_ids).unwrap(); + stmt.bind_text(1, &bucket_ids, Destructor::TRANSIENT)?; stmt } @@ -325,16 +334,17 @@ SELECT match &self.partial { Some(partial) => { // language=SQLite - let updated = self - .db - .prepare_v2( "\ + let updated = self.db.prepare_v2( + "\ UPDATE ps_buckets SET last_applied_op = last_op - WHERE last_applied_op != last_op AND - name IN (SELECT value FROM json_each(json_extract(?1, '$.buckets')))", - )?; - updated.bind_text(1, partial.args, Destructor::STATIC)?; - updated.exec()?; + WHERE last_applied_op != last_op AND name = ?", + )?; + + for bucket in partial.list_buckets() { + updated.bind_text(1, bucket, Destructor::STATIC)?; + updated.exec()?; + } } None => { // language=SQLite From 52f754026cfc6d85cb789e63834e541cdc42bf22 Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Wed, 2 Sep 2026 14:17:15 +0200 Subject: [PATCH 2/3] AI feedback --- crates/core/src/sync/storage_adapter.rs | 2 +- crates/core/src/sync/sync_local.rs | 15 ++++----------- 2 files changed, 5 insertions(+), 12 deletions(-) diff --git a/crates/core/src/sync/storage_adapter.rs b/crates/core/src/sync/storage_adapter.rs index 30084fe..6721a92 100644 --- a/crates/core/src/sync/storage_adapter.rs +++ b/crates/core/src/sync/storage_adapter.rs @@ -281,7 +281,7 @@ WHERE bucket = ?1", self.db, Some(PartialSyncOperation { priority, - checkpoint, + involved_buckets: checkpoint.list_buckets(Some(priority)).collect(), }), now, ), diff --git a/crates/core/src/sync/sync_local.rs b/crates/core/src/sync/sync_local.rs index 19b46f0..bfee930 100644 --- a/crates/core/src/sync/sync_local.rs +++ b/crates/core/src/sync/sync_local.rs @@ -16,7 +16,6 @@ use crate::sync::BucketPriority; use crate::sync::storage_adapter::{ LAST_SEEN_CHECKPOINT_REQUEST_ID_KEY, TARGET_CHECKPOINT_REQUEST_ID_KEY, }; -use crate::sync::streaming_sync::OwnedCheckpoint; use crate::sync::sync_status::TimestampMicros; use crate::utils::SqlBuffer; use crate::utils::database::{Database, Statement}; @@ -26,13 +25,7 @@ use powersync_sqlite_nostd::{self as sqlite, Destructor}; pub struct PartialSyncOperation<'a> { /// The lowest priority part of the partial sync operation. pub priority: BucketPriority, - pub checkpoint: &'a OwnedCheckpoint, -} - -impl<'a> PartialSyncOperation<'a> { - fn list_buckets(&self) -> impl Iterator { - self.checkpoint.list_buckets(Some(self.priority)) - } + pub involved_buckets: Vec<&'a str>, } pub struct SyncOperation<'a> { @@ -321,8 +314,8 @@ SELECT GROUP BY b.row_type, b.row_id;", )?; - let bucket_ids: Vec<&str> = partial.list_buckets().collect(); - let bucket_ids = serde_json::to_string(&bucket_ids).unwrap(); + let bucket_ids = serde_json::to_string(&partial.involved_buckets) + .map_err(PowerSyncError::internal)?; stmt.bind_text(1, &bucket_ids, Destructor::TRANSIENT)?; stmt @@ -341,7 +334,7 @@ SELECT WHERE last_applied_op != last_op AND name = ?", )?; - for bucket in partial.list_buckets() { + for bucket in &partial.involved_buckets { updated.bind_text(1, bucket, Destructor::STATIC)?; updated.exec()?; } From 6e5227b0ebe41a1d50966cb11258286383106b42 Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Tue, 8 Sep 2026 16:40:27 +0200 Subject: [PATCH 3/3] Share logic for prepared statements --- crates/core/src/schema/table_info.rs | 4 +- crates/core/src/sync/sync_local.rs | 221 ++++++++++++--------------- 2 files changed, 101 insertions(+), 124 deletions(-) diff --git a/crates/core/src/schema/table_info.rs b/crates/core/src/schema/table_info.rs index 877ac8b..fc46b9d 100644 --- a/crates/core/src/schema/table_info.rs +++ b/crates/core/src/schema/table_info.rs @@ -336,6 +336,7 @@ impl<'de> Deserialize<'de> for PendingStatement { set.insert(match column { PendingStatementValue::Id => "id".to_string(), PendingStatementValue::Column(name) => name.clone(), + PendingStatementValue::Row => continue, PendingStatementValue::Rest => { rest_parameter_positions.push(i); continue; @@ -366,5 +367,6 @@ pub enum PendingStatementValue { /// Bind to a JSON object containing all columns from the synced row that haven't been matched /// by other statement values. Rest, - // TODO: Stuff like a raw object of put data? + /// The full JSON object for the row, as received from the PowerSync service. + Row, } diff --git a/crates/core/src/sync/sync_local.rs b/crates/core/src/sync/sync_local.rs index bfee930..fb3c449 100644 --- a/crates/core/src/sync/sync_local.rs +++ b/crates/core/src/sync/sync_local.rs @@ -1,8 +1,10 @@ +use core::fmt::Write; + use alloc::collections::btree_map::BTreeMap; -use alloc::format; use alloc::rc::Rc; -use alloc::string::{String, ToString}; +use alloc::string::String; use alloc::vec::Vec; +use alloc::{format, vec}; use serde::Serialize; use serde::ser::SerializeMap; @@ -101,15 +103,6 @@ WHERE target.key = '{TARGET_CHECKPOINT_REQUEST_ID_KEY}' let schema_version = InferredSchemaCache::current_schema_version(self.db)?; let schema_cache = &self.state.inferred_schema_cache; - // We cache the last insert and delete statements for each row - struct CachedStatement { - table: String, - statement: Statement, - } - - let mut last_insert = None::; - let mut last_delete = None::; - let mut untyped_delete_statement: Option = None; let mut untyped_insert_statement: Option = None; @@ -119,10 +112,11 @@ WHERE target.key = '{TARGET_CHECKPOINT_REQUEST_ID_KEY}' let data = statement.column_text(2); if let Some(known) = self.schema.tables.get_mut(type_name) { - if let Some(raw) = &mut known.raw { - match data { - Ok(data) => { - let stmt = raw.put_statement(self.db, schema_version, schema_cache)?; + match data { + Ok(data) => { + let stmt = known.put_statement(self.db, schema_version, schema_cache)?; + + if stmt.needs_parsed_json { 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(|| { @@ -132,70 +126,19 @@ WHERE target.key = '{TARGET_CHECKPOINT_REQUEST_ID_KEY}' })?; let rest = stmt.render_rest_object(json_object)?; - stmt.bind_for_put(id, &json_object, &rest)?; - stmt.exec(type_name, id, Some(&parsed))?; - } - Err(_) => { - let stmt = - raw.delete_statement(self.db, schema_version, schema_cache)?; - stmt.bind_for_delete(id)?; - stmt.exec(type_name, id, None)?; + stmt.bind_for_put(id, data, Some(json_object), rest.as_ref())?; + stmt.exec(type_name, id, Some(&data))?; + } else { + stmt.bind_for_put(id, data, None, None)?; + stmt.exec(type_name, id, Some(&data))?; } } - } else { - // is_err() is essentially a NULL check here. - // NULL data means no PUT operations found, so we delete the row. - if data.is_err() { - // DELETE - let delete_statement = match &last_delete { - Some(stmt) if stmt.table == type_name => &stmt.statement, - _ => { - // Prepare statement when the table changed - let mut statement = SqlBuffer::new(); - statement.push_str("DELETE FROM "); - statement.quote_internal_name(type_name, false); - statement.push_str(" WHERE id = ?"); - - let statement = self.db.prepare_v2(&statement.sql)?; - - &last_delete - .insert(CachedStatement { - table: type_name.to_string(), - statement, - }) - .statement - } - }; - - delete_statement.reset()?; - delete_statement.bind_text(1, id, sqlite::Destructor::STATIC)?; - delete_statement.exec()?; - } else { - // INSERT/UPDATE - let insert_statement = match &last_insert { - Some(stmt) if stmt.table == type_name => &stmt.statement, - _ => { - // Prepare statement when the table changed - let mut statement = SqlBuffer::new(); - statement.push_str("REPLACE INTO "); - statement.quote_internal_name(type_name, false); - statement.push_str("(id, data) VALUES (?, ?)"); - - let statement = self.db.prepare_v2(&statement.sql)?; - - &last_insert - .insert(CachedStatement { - table: type_name.to_string(), - statement, - }) - .statement - } - }; - - insert_statement.reset()?; - insert_statement.bind_text(1, id, sqlite::Destructor::STATIC)?; - insert_statement.bind_text(2, data?, sqlite::Destructor::STATIC)?; - insert_statement.exec()?; + Err(_) => { + // is_err() is essentially a NULL check here. + // NULL data means no PUT operations found, so we delete the row. + let stmt = known.delete_statement(self.db, schema_version, schema_cache)?; + stmt.bind_for_delete(id)?; + stmt.exec(type_name, id, None)?; } } } else { @@ -398,8 +341,10 @@ impl<'a> ParsedDatabaseSchema<'a> { fn add_from_schema(&mut self, schema: &'a Schema) { for raw in &schema.raw_tables { - self.tables - .insert(raw.name.clone(), ParsedSchemaTable::raw(raw)); + self.tables.insert( + raw.name.clone(), + ParsedSchemaTable::new(TableDefinition::Raw(raw)), + ); } } @@ -409,8 +354,12 @@ impl<'a> ParsedDatabaseSchema<'a> { if !table.local_only { let visible_name = table.name; - self.tables - .insert(visible_name, ParsedSchemaTable::json_table()); + self.tables.insert( + visible_name, + ParsedSchemaTable::new(TableDefinition::JsonView { + local_table: table.internal_name, + }), + ); } } @@ -419,25 +368,29 @@ impl<'a> ParsedDatabaseSchema<'a> { } struct ParsedSchemaTable<'a> { - raw: Option>, -} - -struct RawTableWithCachedStatements<'a> { - definition: &'a RawTable, + definition: TableDefinition<'a>, cached_put: Option, cached_delete: Option, } -impl<'a> RawTableWithCachedStatements<'a> { +impl<'a> ParsedSchemaTable<'a> { + const fn new(definition: TableDefinition<'a>) -> Self { + Self { + definition, + cached_put: None, + cached_delete: None, + } + } + fn prepare_lazily( db: Database, slot: &mut Option, - def: Rc, + create_stmt: impl FnOnce() -> Result>, ) -> Result<&PreparedPendingStatement> { Ok(match slot { Some(stmt) => stmt, None => { - let stmt = PreparedPendingStatement::prepare(db, def)?; + let stmt = PreparedPendingStatement::prepare(db, create_stmt()?)?; slot.insert(stmt) } }) @@ -449,14 +402,26 @@ impl<'a> RawTableWithCachedStatements<'a> { schema_version: usize, cache: &InferredSchemaCache, ) -> Result<&'_ PreparedPendingStatement> { - Self::prepare_lazily( - db, - &mut self.cached_put, - match self.definition.put { - Some(ref stmt) => stmt.clone(), - None => cache.infer_put_statement(db, schema_version, &self.definition)?, - }, - ) + Self::prepare_lazily(db, &mut self.cached_put, || { + Ok(match self.definition { + TableDefinition::Raw(raw_table) => match raw_table.put { + Some(ref stmt) => stmt.clone(), + None => cache.infer_put_statement(db, schema_version, raw_table)?, + }, + TableDefinition::JsonView { ref local_table } => { + let mut statement = SqlBuffer::new(); + statement.push_str("REPLACE INTO "); + let _ = statement.identifier().write_str(local_table); + statement.push_str("(id, data) VALUES (?, ?)"); + + Rc::new(PendingStatement { + sql: statement.sql, + params: vec![PendingStatementValue::Id, PendingStatementValue::Row], + named_parameters_index: None, + }) + } + }) + }) } fn delete_statement( @@ -465,36 +430,38 @@ impl<'a> RawTableWithCachedStatements<'a> { schema_version: usize, cache: &InferredSchemaCache, ) -> Result<&'_ PreparedPendingStatement> { - Self::prepare_lazily( - db, - &mut self.cached_delete, - match self.definition.delete { - Some(ref stmt) => stmt.clone(), - None => cache.infer_delete_statement(db, schema_version, &self.definition)?, - }, - ) + Self::prepare_lazily(db, &mut self.cached_delete, || { + Ok(match self.definition { + TableDefinition::Raw(raw_table) => match raw_table.delete { + Some(ref stmt) => stmt.clone(), + None => cache.infer_delete_statement(db, schema_version, raw_table)?, + }, + TableDefinition::JsonView { ref local_table } => { + let mut statement = SqlBuffer::new(); + statement.push_str("DELETE FROM "); + let _ = statement.identifier().write_str(&local_table); + statement.push_str(" WHERE id = ?"); + + Rc::new(PendingStatement { + sql: statement.sql, + params: vec![PendingStatementValue::Id], + named_parameters_index: None, + }) + } + }) + }) } } -impl<'a> ParsedSchemaTable<'a> { - pub const fn json_table() -> Self { - Self { raw: None } - } - - pub fn raw(definition: &'a RawTable) -> Self { - Self { - raw: Some(RawTableWithCachedStatements { - definition, - cached_put: None, - cached_delete: None, - }), - } - } +enum TableDefinition<'a> { + Raw(&'a RawTable), + JsonView { local_table: String }, } struct PreparedPendingStatement { stmt: Statement, definition: Rc, + needs_parsed_json: bool, } impl PreparedPendingStatement { @@ -513,6 +480,10 @@ impl PreparedPendingStatement { Ok(Self { stmt, + needs_parsed_json: pending.params.iter().any(|p| match p { + PendingStatementValue::Id | PendingStatementValue::Row => false, + PendingStatementValue::Column(_) | PendingStatementValue::Rest => true, + }), definition: pending, }) } @@ -565,8 +536,9 @@ impl PreparedPendingStatement { pub fn bind_for_put( &self, id: &str, - json_data: &serde_json::Map, - rest: &Option, + row: &str, + json_data: Option<&serde_json::Map>, + rest: Option<&String>, ) -> Result<()> { use serde_json::Value; @@ -577,8 +549,11 @@ impl PreparedPendingStatement { PendingStatementValue::Id => { self.stmt.bind_text(i, id, Destructor::STATIC)?; } + PendingStatementValue::Row => { + self.stmt.bind_text(i, row, Destructor::STATIC)?; + } PendingStatementValue::Column(column) => { - match json_data.get(column) { + match json_data.and_then(|m| m.get(column)) { Some(Value::Bool(value)) => { self.stmt.bind_int(i, if *value { 1 } else { 0 }) } @@ -634,7 +609,7 @@ impl PreparedPendingStatement { /// Executes the prepared statement, contextualizing errors with the id / data that we've tried /// to insert. - pub fn exec(&self, table: &str, id: &str, data: Option<&serde_json::Value>) -> Result<()> { + pub fn exec(&self, table: &str, id: &str, data: Option<&str>) -> Result<()> { self.stmt.exec().map_err(|e| { let context = match data { None => format!("deleting from {table}, id = {id}"),