Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,3 +1,7 @@
## 0.0.7 (unreleased)

- Update PowerSync core extension to version 0.5.1.

## 0.0.6

- Skip creating `ps_crud` entries when clearing raw tables.
Expand Down
8 changes: 4 additions & 4 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

4 changes: 2 additions & 2 deletions powersync/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -44,8 +44,8 @@ thiserror = "2.0.16"
tokio = { version = "1", features = ["time", "rt"], optional = true }
url = "2.5.7"
serde_with = "3.15.0"
powersync_core = { version = "=0.4.12", features = ["static"] }
powersync_sqlite_nostd = { version = "=0.4.12", features = ["static"] }
powersync_core = { version = "=0.5.1", features = ["static"] }
powersync_sqlite_nostd = { version = "=0.5.1", features = ["static"] }
num-traits = "0.2.19"

[dev-dependencies]
Expand Down
4 changes: 2 additions & 2 deletions powersync/src/db/connection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -48,8 +48,8 @@ impl SqliteConnection {
/// Executes a SQL statement without parameters.
pub fn exec(&self, stmt: &CStr) -> Result<(), PowerSyncError> {
unsafe {
// Safety: We know the stmt is null-terminated.
self.handle().exec(stmt.as_ptr())
// Safety: We're not doing anything that could close the connection.
self.handle().exec(stmt)
}
.map_err(|rc| RawPowerSyncError::RawSqlite {
code: rc,
Expand Down
4 changes: 2 additions & 2 deletions powersync/src/db/core_extension.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,8 +12,8 @@ pub struct CoreExtensionVersion {

impl CoreExtensionVersion {
/// The minimum version of the core extension supported by the native SDK.
pub const MINIMUM: Self = Self::new(0, 4, 7);
pub const MAXIMUM_EXCLUSIVE: Self = Self::new(0, 5, 0);
pub const MINIMUM: Self = Self::new(0, 5, 1);
pub const MAXIMUM_EXCLUSIVE: Self = Self::new(0, 6, 0);

pub const fn new(major: u32, minor: u32, patch: u32) -> Self {
Self {
Expand Down
47 changes: 32 additions & 15 deletions powersync/src/db/internal.rs
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
use crate::db::connection::{SqliteConnection, TransactionGuard, exec_stmt};
use crate::db::connection::{TransactionGuard, exec_stmt};
use crate::schema::SchemaOrCustom;
use crate::{
db::{
Expand All @@ -12,7 +12,7 @@ use crate::{
use event_listener::EventListener;
use futures_lite::future::yield_now;
use futures_lite::{FutureExt, Stream, StreamExt, ready};
use powersync_sqlite_nostd::{Destructor, ResultCode};
use powersync_sqlite_nostd::{ColumnType, Destructor, ResultCode};
use std::sync::{Mutex, Weak};
use std::time::Duration;
use std::{
Expand Down Expand Up @@ -62,28 +62,31 @@ impl InnerPowerSyncState {
let pool = &self.env.pool;
self.did_initialize
.run(|| async {
let conn = pool.writer().await;
let conn = conn.sqlite_connection();
let mut conn = pool.writer().await;
let conn = conn.sqlite_connection_mut();
CoreExtensionVersion::check_from_db(conn)?;

conn.exec(c"SELECT powersync_init()")?;
let tx = TransactionGuard::new(conn)?;
tx.inner.exec(c"SELECT powersync_init()")?;

self.update_schema_internal(conn)?;
self.status.update(|old| old.resolve_offline_state(conn))?;
self.update_schema_internal(&tx)?;
self.status
.update(|old| old.resolve_offline_state(tx.inner))?;
tx.commit()?;

Ok(())
})
.await
.clone()
}

fn update_schema_internal(&self, conn: &SqliteConnection) -> Result<(), PowerSyncError> {
fn update_schema_internal(&self, conn: &TransactionGuard) -> Result<(), PowerSyncError> {
if let SchemaOrCustom::Schema(schema) = self.schema.as_ref() {
schema.validate()?;
};

let serialized_schema = serde_json::to_string(&self.schema)?;
let stmt = conn.prepare("SELECT powersync_replace_schema(?)")?;
let stmt = conn.inner.prepare("SELECT powersync_replace_schema(?)")?;
// Fine because we drop the statement before the serialized schema
stmt.bind_text(1, &serialized_schema, Destructor::STATIC)?;
exec_stmt(stmt)?;
Expand Down Expand Up @@ -118,15 +121,29 @@ impl InnerPowerSyncState {
}
}

Self::set_local_target_op(writer.inner, target_op)?;
Self::target_checkpoint_request_id(&writer, Some(target_op))?;
writer.commit()
}

pub fn set_local_target_op(writer: &SqliteConnection, op: i64) -> Result<(), PowerSyncError> {
let stmt = writer.prepare("UPDATE ps_buckets SET target_op = ? WHERE name = ?")?;
stmt.bind_int64(1, op)?;
stmt.bind_text(2, "$local", Destructor::STATIC)?;
exec_stmt(stmt)
pub fn target_checkpoint_request_id(
writer: &TransactionGuard,
update: Option<i64>,
) -> Result<Option<i64>, PowerSyncError> {
let stmt = writer.inner.prepare("SELECT powersync_control(?, ?);")?;
stmt.bind_text(1, "target_checkpoint_request_id", Destructor::STATIC)?;
if let Some(update) = update {
stmt.bind_int64(2, update)?;
} else {
stmt.bind_null(2)?;
}
let ResultCode::ROW = stmt.step()? else {
panic!("Scalar statement did not return a row")
};

Ok(match stmt.column_type(0)? {
ColumnType::Integer => Some(stmt.column_int64(0)),
_ => None,
})
}

pub async fn reader(&self) -> Result<LeasedConnection, PowerSyncError> {
Expand Down
32 changes: 15 additions & 17 deletions powersync/src/sync/upload.rs
Original file line number Diff line number Diff line change
Expand Up @@ -308,8 +308,10 @@ impl<'a> CrudUpload<'a> {
})
}

fn ps_crud_sequence(conn: &SqliteConnection) -> Result<Option<i64>, PowerSyncError> {
let seq_before = conn.prepare("SELECT seq FROM main.sqlite_sequence WHERE name = ?")?;
fn ps_crud_sequence(tx: &TransactionGuard) -> Result<Option<i64>, PowerSyncError> {
let seq_before = tx
.inner
.prepare("SELECT seq FROM main.sqlite_sequence WHERE name = ?")?;
seq_before.bind_text(1, "ps_crud", Destructor::STATIC)?;

let ResultCode::ROW = seq_before.step()? else {
Expand All @@ -322,21 +324,17 @@ impl<'a> CrudUpload<'a> {
async fn sequence_for_checkpoint(
&self,
) -> Result<Option<PendingCheckpointRequest>, PowerSyncError> {
let reader = self.db.reader().await?;
let reader = reader.sqlite_connection();
{
let stmt =
reader.prepare("SELECT 1 FROM ps_buckets WHERE name = ? AND target_op = ?")?;
stmt.bind_text(1, "$local", Destructor::STATIC)?;
stmt.bind_int64(2, MAX_OP_ID)?;
let mut reader = self.db.reader().await?;
let reader = reader.sqlite_connection_mut();
let read_tx = TransactionGuard::new(reader)?;

let ResultCode::ROW = stmt.step()? else {
// Nothing to update.
return Ok(None);
};
let current_target = InnerPowerSyncState::target_checkpoint_request_id(&read_tx, None)?;
if current_target != Some(MAX_OP_ID) {
// Nothing to update.
return Ok(None);
}

let seq_before = Self::ps_crud_sequence(reader)?;
let seq_before = Self::ps_crud_sequence(&read_tx)?;
Ok(seq_before.map(|seq_before| PendingCheckpointRequest {
crud_sequence: seq_before,
}))
Expand Down Expand Up @@ -369,8 +367,8 @@ impl PendingCheckpointRequest {
return Ok(());
}

let seq_after = CrudUpload::ps_crud_sequence(writer.inner)?
.expect("sqlite sequence should not be empty");
let seq_after =
CrudUpload::ps_crud_sequence(&writer)?.expect("sqlite sequence should not be empty");

if seq_after != self.crud_sequence {
debug!(
Expand All @@ -380,7 +378,7 @@ impl PendingCheckpointRequest {
return Ok(());
}

InnerPowerSyncState::set_local_target_op(writer.inner, op_id)?;
InnerPowerSyncState::target_checkpoint_request_id(&writer, Some(op_id))?;

writer.commit()?;
Ok(())
Expand Down
6 changes: 5 additions & 1 deletion powersync/tests/crud_test.rs
Original file line number Diff line number Diff line change
Expand Up @@ -303,9 +303,13 @@ fn raw_table_clear() {

// Running powersync_clear should delete from users
{
let writer = db.writer().await.unwrap();
let mut writer = db.writer().await.unwrap();
let writer = writer.transaction().unwrap();

let mut stmt = writer.prepare("SELECT powersync_clear(0)").unwrap();
stmt.query_one(params![], |_| Ok(())).unwrap();
drop(stmt);
writer.commit().unwrap();
}

assert_eq!(
Expand Down
Loading