diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index ec4178b..d50c113 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -1,64 +1,54 @@ -name: CI +name: Build and test on: push: branches: - main + pull_request: + branches: + - main permissions: contents: read +env: + CARGO_TERM_COLOR: always + jobs: - build: + rust: + name: Build and test runs-on: ubuntu-latest steps: - - name: Checkout - uses: actions/checkout@v6 + - uses: actions/checkout@v6 with: persist-credentials: false - - name: Install Rust toolchain - uses: dtolnay/rust-toolchain@stable - - - name: Cache cargo - uses: actions/cache@v5 + - name: Cache + uses: actions/cache@v4 with: path: | ~/.cargo/registry ~/.cargo/git target key: ${{ runner.os }}-cargo-${{ hashFiles('**/Cargo.lock') }} + - run: rustup update stable && rustup default stable && rustup component add clippy - - name: Run tests - run: cargo test --all-features + - name: Build all features + run: cargo build --workspace --all-features --verbose - - name: Build - run: cargo build --all-features + - name: Test all features + run: cargo test --workspace --all-features --verbose - lint: - runs-on: ubuntu-latest + macos: + name: Check macOS ARM64 build + runs-on: macos-14 steps: - - name: Checkout - uses: actions/checkout@v6 + - uses: actions/checkout@v6 with: persist-credentials: false - - name: Install Rust toolchain - uses: dtolnay/rust-toolchain@stable - with: - components: rustfmt, clippy - - - name: Cache cargo - uses: actions/cache@v5 - with: - path: | - ~/.cargo/registry - ~/.cargo/git - target - key: ${{ runner.os }}-cargo-${{ hashFiles('**/Cargo.lock') }} - - - name: Check formatting - run: cargo fmt --check + - name: Install stable Rust + run: rustup update stable - - name: Lint with clippy - run: cargo clippy --all-targets --all-features -- -D warnings + - name: Check PowerSync crate + run: cargo check -p powersync --all-targets --all-features diff --git a/.github/workflows/pr.yml b/.github/workflows/pr.yml index 2bcc835..a16aa4f 100644 --- a/.github/workflows/pr.yml +++ b/.github/workflows/pr.yml @@ -34,20 +34,20 @@ jobs: target key: ${{ runner.os }}-cargo-${{ hashFiles('**/Cargo.lock') }} - - name: Format code - run: cargo fmt + - name: Check formatting + run: cargo fmt --all --check - - name: Lint with clippy - run: cargo clippy --all-targets --all-features -- -D warnings + - name: Lint all targets and features + run: cargo clippy --workspace --all-targets --all-features -- -D warnings - - name: Run tests - run: cargo test --all-features + - name: Test all features + run: cargo test --workspace --all-features - - name: Check for security vulnerabilities + - name: Check dependencies run: cargo deny check - - name: Build - run: cargo build --all-features + - name: Build all features + run: cargo build --workspace --all-features osv-scan: uses: google/osv-scanner-action/.github/workflows/osv-scanner-reusable-pr.yml@v2.3.5 diff --git a/.hooks/pre-commit b/.hooks/pre-commit deleted file mode 100755 index 59622e2..0000000 --- a/.hooks/pre-commit +++ /dev/null @@ -1,9 +0,0 @@ -#!/bin/sh -set -e - -echo "==> Running pre-commit checks..." - -echo "==> Checking formatting..." -cargo fmt --check - -echo "==> pre-commit checks passed" diff --git a/.hooks/pre-push b/.hooks/pre-push deleted file mode 100755 index 8c65c55..0000000 --- a/.hooks/pre-push +++ /dev/null @@ -1,12 +0,0 @@ -#!/bin/sh -set -e - -echo "==> Running pre-push checks..." - -echo "==> Running clippy..." -cargo clippy --all-targets --all-features -- -D warnings - -echo "==> Running cargo-deny..." -cargo deny check - -echo "==> pre-push checks passed" diff --git a/CHANGELOG.md b/CHANGELOG.md index 08ed656..c5c0b2a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,6 +1,21 @@ -## 0.0.6 (unreleased) +## v0.0.7-guion.1 (unreleased) + +- Rebase the Guion fork on upstream PowerSync Native v0.0.7. +- Scan queued CRUD immediately when the upload actor connects. +- Release the download writer before awaiting connector or network work. +- Report `ConnectionEstablished` only after a successful HTTP response. +- Accept CRLF framing in JSON sync streams. + +## 0.0.7 + +- Update PowerSync core extension to version 0.5.2. + +## 0.0.6 - Skip creating `ps_crud` entries when clearing raw tables. +- Call `upload_data` repeatedly if an upload fails. +- Add `PowerSyncError::upload_error`, which can be used to convert any error into PowerSync errors for + `upload_data` callbacks. ## 0.0.5 diff --git a/Cargo.lock b/Cargo.lock index 09d6f9b..24a94d9 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -490,9 +490,9 @@ checksum = "c08606f8c3cbf4ce6ec8e28fb0014a2c086708fe954eaa885384a6165172e7e8" [[package]] name = "aws-lc-rs" -version = "1.16.1" +version = "1.18.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "94bffc006df10ac2a68c83692d734a465f8ee6c5b384d8545a636f81d858f4bf" +checksum = "ce2b2dcc879c3bae0d371e77c99f2238400ef24ec001394befa67b6e543add9e" dependencies = [ "aws-lc-sys", "zeroize", @@ -500,14 +500,15 @@ dependencies = [ [[package]] name = "aws-lc-sys" -version = "0.38.0" +version = "0.44.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4321e568ed89bb5a7d291a7f37997c2c0df89809d7b6d12062c81ddb54aa782e" +checksum = "f09fae7be8bb3174e05c6afdb34199e6dc0c7c04ba9fa237b1967adfbde27483" dependencies = [ "cc", "cmake", "dunce", "fs_extra", + "pkg-config", ] [[package]] @@ -1287,11 +1288,10 @@ checksum = "dea2df4cf52843e0452895c455a1a2cfbb842a1e7329671acf418fdc53ed4c59" [[package]] name = "event-listener" -version = "5.4.1" +version = "5.4.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e13b66accf52311f30a0db42147dadea9850cb48cd070028831ae5f5d4b856ab" +checksum = "5a23add41df1562121a9393cb065eab5146a1242410f23a644851e90cfd669d2" dependencies = [ - "concurrent-queue", "parking", "pin-project-lite", ] @@ -1777,15 +1777,6 @@ version = "0.12.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8a9ee70c43aaf417c914396645a0fa852624801b24ebb7ae78fe8272889ac888" -[[package]] -name = "hashbrown" -version = "0.14.5" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e5274423e17b7c9fc20b6e7e208532f9b19825d82dfd615708b70edd83df41f1" -dependencies = [ - "ahash", -] - [[package]] name = "hashbrown" version = "0.15.5" @@ -1806,11 +1797,11 @@ dependencies = [ [[package]] name = "hashlink" -version = "0.9.1" +version = "0.11.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6ba4ff7128dee98c7dc9794b6a411377e1404dba1c97deb8d1a55297bd25d8af" +checksum = "ea0b22561a9c04a7cb1a302c013e0259cd3b4bb619f145b32f72b8b4bcbed230" dependencies = [ - "hashbrown 0.14.5", + "hashbrown 0.16.1", ] [[package]] @@ -2325,9 +2316,9 @@ dependencies = [ [[package]] name = "libsqlite3-sys" -version = "0.30.1" +version = "0.37.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2e99fb7a497b1e3339bc746195567ed8d3e24945ecd636e3619d20b9de9e9149" +checksum = "b1f111c8c41e7c61a49cd34e44c7619462967221a6443b0ec299e0ac30cfb9b1" dependencies = [ "cc", "pkg-config", @@ -3093,7 +3084,7 @@ checksum = "439ee305def115ba05938db6eb1644ff94165c5ab5e9420d1c1bcedbba909391" [[package]] name = "powersync" -version = "0.0.5" +version = "0.0.7" dependencies = [ "async-channel", "async-executor", @@ -3108,8 +3099,10 @@ dependencies = [ "futures-lite", "futures-test", "log", + "num-traits", "pin-project-lite", "powersync_core", + "powersync_sqlite_nostd", "powersync_test_utils", "reqwest", "rusqlite", @@ -3124,9 +3117,9 @@ dependencies = [ [[package]] name = "powersync_core" -version = "0.4.12" +version = "0.5.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "77f9c7f0117bda7f68ca872528e0c2c78de94d34f81a6e167e2046e8578bd080" +checksum = "32ab437bc28b8007e789d4d4b39861a7dbd1eca102d4c453f3ba12d320b0de7a" dependencies = [ "bytes", "const_format", @@ -3144,9 +3137,9 @@ dependencies = [ [[package]] name = "powersync_sqlite_nostd" -version = "0.4.12" +version = "0.5.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "41a5a95198e2ab901965138fced0315621073ef557a04b5552a51774658abf13" +checksum = "f4e318c28daeba1a83f93eb42917972b243a5def587eb27a7c7e5da938df3b67" dependencies = [ "bindgen", "num-derive 0.4.2", @@ -3556,11 +3549,21 @@ dependencies = [ "windows-sys 0.52.0", ] +[[package]] +name = "rsqlite-vfs" +version = "0.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a8a1f2315036ef6b1fbacd1972e8ee7688030b0a2121edfc2a6550febd41574d" +dependencies = [ + "hashbrown 0.16.1", + "thiserror 2.0.18", +] + [[package]] name = "rusqlite" -version = "0.32.1" +version = "0.39.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7753b721174eb8ff87a9a0e799e2d7bc3749323e773db92e0984debb00019d6e" +checksum = "a0d2b0146dd9661bf67bb107c0bb2a55064d556eeb3fc314151b957f313bcd4e" dependencies = [ "bitflags 2.11.0", "fallible-iterator", @@ -3568,6 +3571,7 @@ dependencies = [ "hashlink", "libsqlite3-sys", "smallvec", + "sqlite-wasm-rs", ] [[package]] @@ -4029,6 +4033,18 @@ dependencies = [ "bitflags 2.11.0", ] +[[package]] +name = "sqlite-wasm-rs" +version = "0.5.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2f4206ed3a67690b9c29b77d728f6acc3ce78f16bf846d83c94f76400320181b" +dependencies = [ + "cc", + "js-sys", + "rsqlite-vfs", + "wasm-bindgen", +] + [[package]] name = "stable_deref_trait" version = "1.2.1" diff --git a/README.md b/README.md index 78b4812..b6f33d3 100644 --- a/README.md +++ b/README.md @@ -13,6 +13,9 @@ _[PowerSync](https://www.powersync.com) is a sync engine for building local-firs This repository contains code used to build a PowerSync SDK for native development. PowerSync is available as a Rust crate in `powersync/`, and on crates.io as the `powersync` crate. +Guion release tags add a small, tested compatibility layer for Guion consumers. See +[`docs/guion-patches.md`](docs/guion-patches.md) for the exact upstream delta and patch policy. + ## Running the examples To start an example: @@ -23,7 +26,7 @@ To start an example: 3. Compile and run an example here: `cargo run -p egui_todolist`. ```yaml -# Sync-rule docs: https://docs.powersync.com/usage/sync-rules +# Sync Streams docs: https://docs.powersync.com/sync/streams/overview streams: lists: query: SELECT * FROM lists #WHERE owner_id = auth.user_id() @@ -32,5 +35,5 @@ streams: query: SELECT * FROM todos WHERE list_id = subscription.parameter('list') #AND list_id IN (SELECT id FROM lists WHERE owner_id = auth.user_id()) config: - edition: 2 + edition: 3 ``` diff --git a/deny.toml b/deny.toml index a1b5dfb..2bc230f 100644 --- a/deny.toml +++ b/deny.toml @@ -20,22 +20,22 @@ allow = [ "OpenSSL", "Ubuntu-font-1.0", ] -exceptions = [] -# Workspace crates are private/unlicensed + [[licenses.clarify]] crate = "powersync" expression = "Apache-2.0" license-files = [] + [[licenses.clarify]] crate = "powersync_test_utils" expression = "Apache-2.0" license-files = [] -# egui_todolist is an example app — treat as Apache-2.0 + [[licenses.clarify]] crate = "egui_todolist" expression = "Apache-2.0" license-files = [] -# aws-lc-sys pulls in OpenSSL — add exception for ISC+OpenSSL combined license + [[licenses.clarify]] crate = "aws-lc-sys" expression = "ISC AND (Apache-2.0 OR ISC) AND OpenSSL" @@ -47,16 +47,11 @@ wildcards = "allow" [advisories] ignore = [ - # aws-lc-sys transitive advisory — upstream is aws-lc - "RUSTSEC-2026-0047", # AWS-LC: CN wildcard bypass - "RUSTSEC-2026-0048", # AWS-LC: CRL DP scope check error - "RUSTSEC-2026-0044", # AWS-LC: CN wildcard bypass via Unicode - # rustls-webpki advisory — transitive via reqwest's rustls - "RUSTSEC-2026-0049", # rustls-webpki: CRL matching logic error - # remove_dir_all — transitive via powersync_test_utils dev-dependency tempdir - "RUSTSEC-2023-0018", # remove_dir_all: TOCTOU race condition - # tempdir unmaintained — transitive via powersync_test_utils - "RUSTSEC-2018-0017", # tempdir deprecated + "RUSTSEC-2023-0018", # remove_dir_all is a test-only tempdir dependency + "RUSTSEC-2018-0017", # tempdir is test-only and inherited from upstream + "RUSTSEC-2026-0195", # quick-xml NsReader is build-only via Wayland scanner on trusted XML + "RUSTSEC-2026-0194", # quick-xml is used only by upstream example desktop build tooling + "RUSTSEC-2026-0192", # ttf-parser is used only by the upstream egui example ] [sources] diff --git a/docs/guion-patches.md b/docs/guion-patches.md new file mode 100644 index 0000000..aed0c37 --- /dev/null +++ b/docs/guion-patches.md @@ -0,0 +1,32 @@ +# Guion fork patch inventory + +The Guion fork is rebuilt from upstream PowerSync Native v0.0.7. A patch stays in a Guion +release only when its current failure mode is covered by a deterministic regression test or it is +an explicit Guion build policy. + +## Retained delta + +| Patch | Failure mode | Regression evidence | Upstream v0.0.7 | Disposition | +| --- | --- | --- | --- | --- | +| Connect-time CRUD scan | CRUD queued before `connect()` can miss the early notification and remain stranded | `connect_uploads_crud_that_was_already_queued` | Missing | Retain; upstream candidate | +| Download writer scope | Connector credential/network awaits can retain the only writer and deadlock other writes | `fetching_credentials_does_not_hold_the_download_writer_lease` | Missing | Retain; upstream candidate | +| Connection status ordering | Transport and non-2xx errors can emit `ConnectionEstablished` before the error | `sync::download::http::tests` | Missing | Retain; upstream candidate | +| CRLF framing | JSON lines ending in CRLF expose a trailing `\r` to the parser | `util::line_split::test` | Missing | Retain; upstream candidate | + +Upstream v0.0.7 already retries a failed `upload_data` call in the same upload cycle. Its +`upload_retry` test remains the source of truth; the fork does not add another retry worker. + +## Dropped legacy patches + +| Legacy patch | Why it is absent from v0.0.7-guion.1 | +| --- | --- | +| rusqlite 0.32 alignment and API adaptations | The SQLx SQLite consumer is being removed; use upstream optional rusqlite 0.39 and `powersync_sqlite_nostd` 0.5.2. | +| Reader `busy_timeout` | No deterministic failure remains with one PowerSync-owned pool. The upstream writer keeps its 30-second timeout. | +| Extra `BEGIN IMMEDIATE` sites | The PowerSync writer mutex serializes SDK writes; no current `BUSY_SNAPSHOT` regression justifies broader locking. | +| Reader lease release sender | The upstream lease can only be constructed when `PoolReaders` exists and returns through the same shared pool state. | +| `From` for `PowerSyncError` | No SDK or Guion consumer path uses this public conversion. | +| Broad Clippy allows | The private-interface warning is fixed by narrowing internal subscription command visibility. | + +Do not restore a dropped patch based on suspicion or an intermittent stress failure. First add a +minimal deterministic test that identifies the failing invariant, then retain only the smallest +fix for that test. diff --git a/examples/egui_todolist/Cargo.toml b/examples/egui_todolist/Cargo.toml index ed3635a..c6eb127 100644 --- a/examples/egui_todolist/Cargo.toml +++ b/examples/egui_todolist/Cargo.toml @@ -11,7 +11,7 @@ futures-lite = "2.6.1" log = "0.4.28" powersync = { path = "../../powersync", features = ["tokio", "reqwest"] } reqwest = { version = "0.13.2", features = ["json"] } -rusqlite = { version = "0.32.0", features = ["load_extension", "bundled"] } +rusqlite = { version = "0.39.0", features = ["load_extension", "bundled"] } serde = "1.0.228" serde_json = "1.0.145" tokio = { version = "1.47.1", features = ["rt-multi-thread", "net"] } diff --git a/lefthook.yml b/lefthook.yml index f29206f..3157a57 100644 --- a/lefthook.yml +++ b/lefthook.yml @@ -2,12 +2,12 @@ pre-commit: parallel: true commands: cargo-fmt: - run: cargo fmt --check + run: cargo fmt --all --check pre-push: parallel: true commands: cargo-clippy: - run: cargo clippy --all-targets --all-features -- -D warnings + run: cargo clippy --workspace --all-targets --all-features -- -D warnings cargo-deny: run: cargo deny check diff --git a/powersync/Cargo.toml b/powersync/Cargo.toml index cc4a9b4..33ddfa7 100644 --- a/powersync/Cargo.toml +++ b/powersync/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "powersync" -version = "0.0.5" +version = "0.0.7" edition = "2024" license = "Apache-2.0" @@ -15,9 +15,12 @@ description = "Experimental PowerSync SDK for Rust applications." crate-type = ["lib"] [features] +default = ["rusqlite"] + tokio = ["dep:tokio"] smol = ["dep:async-io"] reqwest = ["dep:reqwest"] +rusqlite = ["dep:rusqlite"] ffi = [] [dependencies] @@ -29,11 +32,11 @@ async-oneshot = "0.5.9" atomic_enum = "0.3.0" event-listener = "5.4.1" futures-lite = "2.6.1" -reqwest = { version = "0.13.2", default-features = false, optional = true, features = ["stream", "rustls"] } +reqwest = { version = "0.13.2", optional = true, features = ["stream"] } bytes = "1" log = "0.4.28" pin-project-lite = "0.2.16" -rusqlite = { version = "0.32.0", features = ["load_extension"] } +rusqlite = { version = "0.39.0", optional = true, features = ["load_extension"] } scopeguard = "1.2.0" serde = { version = "1.0.219", features = ["derive", "rc"] } serde_json = { version = "1.0.143", features = ["raw_value"] } @@ -41,7 +44,9 @@ 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_core = { version = "=0.5.2", features = ["static"] } +powersync_sqlite_nostd = { version = "=0.5.2", features = ["static"] } +num-traits = "0.2.19" [dev-dependencies] async-executor = "1.13.3" @@ -49,7 +54,3 @@ async-task = "4.7.1" futures-lite = "2.6.1" futures-test = "0.3.31" powersync_test_utils = { path = "../powersync_test_utils" } - -[lints.clippy] -clone_on_copy = "allow" -derivable_impls = "allow" diff --git a/powersync/src/db/connection.rs b/powersync/src/db/connection.rs new file mode 100644 index 0000000..849da56 --- /dev/null +++ b/powersync/src/db/connection.rs @@ -0,0 +1,317 @@ +use crate::error::{PowerSyncError, RawPowerSyncError}; +use num_traits::cast::FromPrimitive; +use powersync_sqlite_nostd::bindings::{sqlite3_close_v2, sqlite3_open_v2}; +use powersync_sqlite_nostd::{Connection, ManagedConnection, ManagedStmt, ResultCode, sqlite3}; +use std::ffi::{CStr, CString, c_int}; +use std::path::Path; +use std::ptr::{null, null_mut}; + +/// The SQLite connection used by the PowerSync Rust SDK. +/// +/// When the `rusqlite` feature is enabled, we use rusqlite connections. +/// Without that feature, we use raw `*mut sqlite3` pointers. Disabling that +/// feature can be useful when a custom SQLite build (e.g. `sqlite3mc`) needs +/// to be used with the SDK. +pub struct SqliteConnection { + #[cfg(not(feature = "rusqlite"))] + raw: RawSqliteConnection, + #[cfg(feature = "rusqlite")] + inner: rusqlite::Connection, +} + +impl SqliteConnection { + /// Returns the `*mut sqlite3` pointer from the inner connection. + /// + /// This method is unsafe since the pointer could be used to transform the connection + /// into an unexpected state. + #[cfg(feature = "rusqlite")] + pub unsafe fn handle(&self) -> *mut sqlite3 { + unsafe { self.inner.handle() }.cast() + } + + #[cfg(not(feature = "rusqlite"))] + pub unsafe fn handle(&self) -> *mut sqlite3 { + self.raw.0.db + } + + #[cfg(feature = "rusqlite")] + pub fn rusqlite_connection(&self) -> &rusqlite::Connection { + &self.inner + } + + #[cfg(feature = "rusqlite")] + pub fn rusqlite_connection_mut(&mut self) -> &mut rusqlite::Connection { + &mut self.inner + } + + /// Executes a SQL statement without parameters. + pub fn exec(&self, stmt: &CStr) -> Result<(), PowerSyncError> { + unsafe { + // Safety: We're not doing anything that could close the connection. + self.handle().exec(stmt) + } + .map_err(|rc| RawPowerSyncError::RawSqlite { + code: rc, + context: format!("Could not run {}", stmt.to_string_lossy()), + })?; + + Ok(()) + } + + pub fn prepare(&self, stmt: &str) -> Result { + unsafe { + // Safety: We're not doing anything that could close the connection. + self.handle() + } + .prepare_v2(stmt) + .map_err(|rc| { + RawPowerSyncError::RawSqlite { + code: rc, + context: format!("Could not prepare {stmt}"), + } + .into() + }) + } +} + +/// Utility for running a block in a transaction. +pub struct TransactionGuard<'a> { + pub inner: &'a mut SqliteConnection, + active: bool, +} + +impl<'a> TransactionGuard<'a> { + pub fn new(connection: &'a mut SqliteConnection) -> Result { + if !unsafe { connection.handle().get_autocommit() } { + return Err(PowerSyncError::argument_error( + "Connection already in transaction", + )); + } + + connection.exec(c"BEGIN")?; + Ok(TransactionGuard { + inner: connection, + active: true, + }) + } + + pub fn commit(mut self) -> Result<(), PowerSyncError> { + self.inner.exec(c"COMMIT")?; + self.active = false; + Ok(()) + } + + fn rollback_internal(&mut self) -> Result<(), PowerSyncError> { + self.inner.exec(c"ROLLBACK") + } +} + +impl Drop for TransactionGuard<'_> { + fn drop(&mut self) { + if self.active { + // Rollback if the transaction hasn't explicitly been committed. + let _ = self.rollback_internal(); + } + } +} + +#[cfg(feature = "rusqlite")] +impl From for SqliteConnection { + fn from(value: rusqlite::Connection) -> Self { + Self { inner: value } + } +} + +#[cfg(not(feature = "rusqlite"))] +impl From for SqliteConnection { + fn from(value: RawSqliteConnection) -> Self { + Self { raw: value } + } +} + +#[cfg(feature = "rusqlite")] +impl From for SqliteConnection { + fn from(value: RawSqliteConnection) -> Self { + let conn = value.0.db; + + // Don't call sqlite3_close_v2, we want to transfer ownership. + let _ = std::mem::ManuallyDrop::new(value.0); + + Self { + inner: unsafe { + // Safety: The never dropped ManuallyDrop transfers ownership from the + // RawSqliteConnection to rusqlite. + rusqlite::Connection::from_handle_owned(conn.cast()) + } + .unwrap(), + } + } +} + +pub struct RawSqliteConnection(ManagedConnection); + +unsafe impl Send for RawSqliteConnection {} + +impl RawSqliteConnection { + pub fn open(path: &CStr, flags: u32) -> Result { + let mut db = null_mut(); + let rc = ResultCode::from_i32(unsafe { + sqlite3_open_v2(path.as_ptr(), &mut db, flags as c_int, null()) + }) + .unwrap(); + + if rc == ResultCode::OK { + Ok(Self(ManagedConnection { db })) + } else { + if !db.is_null() { + // SQLite may allocate an error-bearing handle even when open fails. + // No statements can exist yet, so closing it here releases all resources. + unsafe { + sqlite3_close_v2(db); + } + } + Err(RawPowerSyncError::RawSqlite { + code: rc, + context: format!("Could not open database {}", path.to_string_lossy()), + } + .into()) + } + } + + pub fn open_path>(path: P, flags: u32) -> Result { + Self::open(path_to_cstring(path.as_ref())?.as_ref(), flags) + } +} + +pub fn exec_stmt(stmt: ManagedStmt) -> Result<(), PowerSyncError> { + while let ResultCode::ROW = stmt.step().map_err(|e| RawPowerSyncError::RawSqlite { + code: e, + context: format!("Stepping through {}", stmt.sql().unwrap_or("unknown SQL")), + })? { + // Keep stepping through statement. + } + + Ok(()) +} + +#[cfg(unix)] +fn path_to_cstring(p: &Path) -> Result { + use std::os::unix::ffi::OsStrExt; + Ok( + CString::new(p.as_os_str().as_bytes()).map_err(|_| RawPowerSyncError::ArgumentError { + desc: format!("Invalid path: {p:?}").into(), + })?, + ) +} + +#[cfg(all(test, feature = "rusqlite"))] +mod tests { + use std::process::Command; + use std::sync::atomic::{AtomicUsize, Ordering}; + + use powersync_sqlite_nostd::bindings::{ + SQLITE_OPEN_CREATE, SQLITE_OPEN_READWRITE, sqlite3_memory_used, + }; + + use super::*; + + static NEXT_TEST_DATABASE: AtomicUsize = AtomicUsize::new(0); + const SQLITE_MEMORY_TEST: &str = + "db::connection::tests::repeated_open_failures_do_not_leak_sqlite_handles"; + const SQLITE_MEMORY_TEST_CHILD: &str = "POWERSYNC_SQLITE_MEMORY_TEST_CHILD"; + + fn test_database_path(name: &str) -> std::path::PathBuf { + std::env::temp_dir().join(format!( + "powersync-{name}-{}-{}.sqlite", + std::process::id(), + NEXT_TEST_DATABASE.fetch_add(1, Ordering::Relaxed) + )) + } + + #[test] + fn failed_commit_rolls_back_before_returning_connection() { + let path = test_database_path("commit-rollback"); + let setup = rusqlite::Connection::open(&path).unwrap(); + setup + .execute_batch( + "PRAGMA journal_mode = DELETE; + CREATE TABLE values_table (value INTEGER NOT NULL); + INSERT INTO values_table VALUES (1);", + ) + .unwrap(); + + let reader = rusqlite::Connection::open(&path).unwrap(); + reader.execute_batch("BEGIN").unwrap(); + let _: i64 = reader + .query_row("SELECT value FROM values_table", [], |row| row.get(0)) + .unwrap(); + + let mut writer = SqliteConnection::from(rusqlite::Connection::open(&path).unwrap()); + let tx = TransactionGuard::new(&mut writer).unwrap(); + tx.inner.exec(c"UPDATE values_table SET value = 2").unwrap(); + + assert!(tx.commit().is_err()); + assert!(unsafe { writer.handle().get_autocommit() }); + + reader.execute_batch("ROLLBACK").unwrap(); + let value: i64 = writer + .rusqlite_connection() + .query_row("SELECT value FROM values_table", [], |row| row.get(0)) + .unwrap(); + assert_eq!(value, 1); + + drop(reader); + drop(setup); + drop(writer); + std::fs::remove_file(path).unwrap(); + } + + #[test] + fn repeated_open_failures_do_not_leak_sqlite_handles() { + if std::env::var_os(SQLITE_MEMORY_TEST_CHILD).is_none() { + let status = Command::new(std::env::current_exe().unwrap()) + .arg(SQLITE_MEMORY_TEST) + .arg("--exact") + .env(SQLITE_MEMORY_TEST_CHILD, "1") + .status() + .unwrap(); + assert!( + status.success(), + "SQLite memory test subprocess failed: {status}" + ); + return; + } + + let path = test_database_path("missing-parent").join("database.sqlite"); + let flags = SQLITE_OPEN_READWRITE | SQLITE_OPEN_CREATE; + + // Warm SQLite's process-global caches before measuring per-open resources. + for _ in 0..128 { + assert!(RawSqliteConnection::open_path(&path, flags).is_err()); + } + let memory_before = unsafe { sqlite3_memory_used() }; + + for _ in 0..128 { + assert!(RawSqliteConnection::open_path(&path, flags).is_err()); + } + + let memory_after = unsafe { sqlite3_memory_used() }; + assert!( + memory_after - memory_before < 1024, + "failed opens leaked {} SQLite bytes", + memory_after - memory_before + ); + } +} + +#[cfg(not(unix))] +fn path_to_cstring(p: &Path) -> Result { + let s = p.to_str().ok_or_else(|| RawPowerSyncError::ArgumentError { + desc: format!("Invalid path: {p:?}").into(), + })?; + Ok( + CString::new(s).map_err(|_| RawPowerSyncError::ArgumentError { + desc: format!("Invalid path: {p:?}").into(), + })?, + ) +} diff --git a/powersync/src/db/core_extension.rs b/powersync/src/db/core_extension.rs index ccc9a19..f35e66e 100644 --- a/powersync/src/db/core_extension.rs +++ b/powersync/src/db/core_extension.rs @@ -1,8 +1,7 @@ -use std::{fmt::Display, str::FromStr}; - -use rusqlite::{Connection, params}; - +use crate::db::connection::SqliteConnection; use crate::error::{PowerSyncError, RawPowerSyncError}; +use powersync_sqlite_nostd::ResultCode; +use std::{fmt::Display, str::FromStr}; #[derive(Clone, PartialEq, PartialOrd, Eq, Ord)] pub struct CoreExtensionVersion { @@ -13,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 { @@ -35,17 +34,13 @@ impl CoreExtensionVersion { } } - pub(crate) fn check_from_db(conn: &Connection) -> Result { - let version = - conn.prepare("SELECT powersync_rs_version()")? - .query_row(params![], |row| { - let value = row.get_ref(0)?; - value - .as_str()? - .parse::() - .map_err(|_| rusqlite::Error::InvalidQuery) - })?; + pub(crate) fn check_from_db(conn: &SqliteConnection) -> Result { + let stmt = conn.prepare("SELECT powersync_rs_version()")?; + let ResultCode::ROW = stmt.step()? else { + panic!("Expected row") // Can't happen, scalar select + }; + let version = stmt.column_text(0)?.parse::()?; version.validate()?; Ok(version) } diff --git a/powersync/src/db/crud.rs b/powersync/src/db/crud.rs index 8a3e7b2..763c4bd 100644 --- a/powersync/src/db/crud.rs +++ b/powersync/src/db/crud.rs @@ -3,12 +3,12 @@ use std::task::{Context, Poll}; use futures_lite::{FutureExt, Stream, ready}; use pin_project_lite::pin_project; -use rusqlite::params; +use powersync_sqlite_nostd::ResultCode; use serde::{Deserialize, Serialize}; use serde_json::{Map, Value}; use crate::PowerSyncDatabase; -use crate::error::{PowerSyncError, RawPowerSyncError}; +use crate::error::PowerSyncError; /// All local writes that were made in a single SQLite transaction. pub struct CrudTransaction<'a> { @@ -147,17 +147,19 @@ impl<'a> CrudTransactionStream<'a> { ) -> Result)>, PowerSyncError> { let last = last.unwrap_or(-1); let reader = db.reader().await?; - let mut stmt = reader.prepare_cached(Self::SQL)?; - let mut rows = stmt.query(params![last])?; + let conn = reader.sqlite_connection(); + let stmt = conn.prepare(Self::SQL)?; + stmt.bind_int64(1, last)?; + let mut crud_entries = vec![]; let mut last = None::<(i64, i64)>; - while let Some(row) = rows.next()? { - let id: i64 = row.get(0)?; - let tx_id: i64 = row.get(1)?; - let data = row.get_ref(2)?.as_str().map_err(RawPowerSyncError::from)?; - last = Some((id, tx_id)); + while let ResultCode::ROW = stmt.step()? { + let id = stmt.column_int64(0); + let tx_id = stmt.column_int64(1); + let data = stmt.column_text(2)?; + last = Some((id, tx_id)); crud_entries.push(CrudEntry::parse(id, tx_id, data)?); } diff --git a/powersync/src/db/internal.rs b/powersync/src/db/internal.rs index 0be998e..2eafeab 100644 --- a/powersync/src/db/internal.rs +++ b/powersync/src/db/internal.rs @@ -1,14 +1,4 @@ -use event_listener::EventListener; -use futures_lite::{FutureExt, Stream, StreamExt, ready}; -use rusqlite::{Connection, TransactionBehavior, params}; -use std::sync::{Mutex, Weak}; -use std::time::Duration; -use std::{ - pin::Pin, - sync::Arc, - task::{Context, Poll}, -}; - +use crate::db::connection::{TransactionGuard, exec_stmt}; use crate::schema::SchemaOrCustom; use crate::{ db::{ @@ -19,6 +9,17 @@ use crate::{ sync::{MAX_OP_ID, coordinator::SyncCoordinator, status::SyncStatus, status::SyncStatusData}, util::SharedFuture, }; +use event_listener::EventListener; +use futures_lite::future::yield_now; +use futures_lite::{FutureExt, Stream, StreamExt, ready}; +use powersync_sqlite_nostd::{ColumnType, Destructor, ResultCode}; +use std::sync::{Mutex, Weak}; +use std::time::Duration; +use std::{ + pin::Pin, + sync::Arc, + task::{Context, Poll}, +}; pub struct InnerPowerSyncState { /// External dependencies (timers and HTTP clients) used to implement SDK functionality. @@ -61,14 +62,17 @@ impl InnerPowerSyncState { let pool = &self.env.pool; self.did_initialize .run(|| async { - let conn = pool.writer().await; - CoreExtensionVersion::check_from_db(&conn)?; + let mut conn = pool.writer().await; + let conn = conn.sqlite_connection_mut(); + CoreExtensionVersion::check_from_db(conn)?; - conn.prepare("SELECT powersync_init()")? - .query_row(params![], |_| Ok(()))?; + 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(()) }) @@ -76,14 +80,17 @@ impl InnerPowerSyncState { .clone() } - fn update_schema_internal(&self, conn: &Connection) -> 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)?; - conn.prepare("SELECT powersync_replace_schema(?)")? - .query_row(params![serialized_schema], |_| Ok(()))?; + 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)?; + // TODO: Update readers? Should be fine at the moment because we're only doing this during // initialization. Ok(()) @@ -97,25 +104,46 @@ impl InnerPowerSyncState { write_checkpoint: Option, ) -> Result<(), PowerSyncError> { let mut writer = self.writer().await?; - let writer = writer.transaction_with_behavior(TransactionBehavior::Immediate)?; + let writer = TransactionGuard::new(writer.sqlite_connection_mut())?; + + { + let stmt = writer.inner.prepare("DELETE FROM ps_crud WHERE id <= ?")?; + stmt.bind_int64(1, last_client_id)?; + exec_stmt(stmt)?; + } - writer.execute("DELETE FROM ps_crud WHERE id <= ?", params![last_client_id])?; let mut target_op: i64 = MAX_OP_ID; if let Some(write_checkpoint) = write_checkpoint { // If there are no remaining crud items we can set the target op to the checkpoint. - let mut stmt = writer.prepare("SELECT 1 FROM ps_crud LIMIT 1")?; - if stmt.query(params![])?.next()?.is_none() { + let stmt = writer.inner.prepare("SELECT 1 FROM ps_crud LIMIT 1")?; + if let ResultCode::DONE = stmt.step()? { target_op = write_checkpoint; } } - writer.execute( - "UPDATE ps_buckets SET target_op = ? WHERE name = ?", - params![target_op, "$local"], - )?; - writer.commit()?; + Self::target_checkpoint_request_id(&writer, Some(target_op))?; + writer.commit() + } - Ok(()) + pub fn target_checkpoint_request_id( + writer: &TransactionGuard, + update: Option, + ) -> Result, 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 { @@ -134,8 +162,12 @@ impl InnerPowerSyncState { *guard }; - if let Some(delay) = delay { + if let Some(delay) = delay + && delay > Duration::ZERO + { self.env.timer.delay_once(delay).await + } else { + yield_now().await } } diff --git a/powersync/src/db/mod.rs b/powersync/src/db/mod.rs index 4f843a3..bae02a9 100644 --- a/powersync/src/db/mod.rs +++ b/powersync/src/db/mod.rs @@ -17,11 +17,10 @@ use crate::{ error::PowerSyncError, sync::{download::DownloadActor, status::SyncStatusData, upload::UploadActor}, }; -use futures_lite::stream::{once, once_future}; use futures_lite::{FutureExt, Stream, StreamExt}; -use rusqlite::{Params, Statement, params}; mod async_support; +pub(crate) mod connection; pub mod core_extension; pub mod crud; pub(crate) mod internal; @@ -147,14 +146,16 @@ impl PowerSyncDatabase { /// This method is a core building block for reactive applications with PowerSync - since it /// updates automatically, all writes (regardless of whether they're local or due to synced /// writes from your backend) are reflected. - pub fn watch_statement( + #[cfg(feature = "rusqlite")] + pub fn watch_statement( &self, sql: String, params: P, read: F, ) -> impl Stream> + 'static where - for<'a> F: (Fn(&'a mut Statement, P) -> Result) + 'static + Clone, + for<'a> F: + (Fn(&'a mut rusqlite::Statement, P) -> Result) + 'static + Clone, { // Find and watch referenced tables. We assume the set of read tables is fixed for a given // SQL query and parameters. We also want this to emit initially without an update so that @@ -186,14 +187,15 @@ impl PowerSyncDatabase { }) } + #[cfg(feature = "rusqlite")] fn emit_on_statement_changes( &self, emit_initially: bool, sql: String, - params: impl Params + 'static, + params: impl rusqlite::Params + 'static, ) -> impl Stream> + 'static { // Stream emitting referenced tables once. - let tables = once_future(self.clone().find_tables(sql, params)); + let tables = futures_lite::stream::once_future(self.clone().find_tables(sql, params)); // Stream emitting updates, or a single error if we couldn't resolve tables. let db = self.clone(); @@ -202,14 +204,15 @@ impl PowerSyncDatabase { .watch_tables(emit_initially, referenced_tables) .map(Ok) .boxed(), - Err(e) => once(Err(e)).boxed(), + Err(e) => futures_lite::stream::once(Err(e)).boxed(), }) } /// Finds all tables that are used in a given select statement. /// /// This can be used together with [watch_tables] to build an auto-updating stream of queries. - async fn find_tables( + #[cfg(feature = "rusqlite")] + async fn find_tables( self, sql: impl Into>, params: P, @@ -231,7 +234,7 @@ impl PowerSyncDatabase { && matches!(p3.as_i64(), Ok(0)) && let Ok(page) = p2.as_i64() { - let mut found_table = find_table_stmt.query(params![page])?; + let mut found_table = find_table_stmt.query(rusqlite::params![page])?; if let Some(found_table) = found_table.next()? { let table_name: String = found_table.get(0)?; found_tables.insert(table_name); diff --git a/powersync/src/db/pool.rs b/powersync/src/db/pool.rs index 9014b0e..763c51f 100644 --- a/powersync/src/db/pool.rs +++ b/powersync/src/db/pool.rs @@ -1,16 +1,16 @@ -use std::{ - collections::HashSet, - mem::MaybeUninit, - ops::{Deref, DerefMut}, - path::Path, - sync::Arc, -}; +#[cfg(feature = "rusqlite")] +use std::ops::{Deref, DerefMut}; +use std::{collections::HashSet, mem::MaybeUninit, path::Path, sync::Arc}; use async_channel::{Receiver, Sender}; use async_lock::{Mutex, MutexGuardArc}; -use rusqlite::{Connection, Error, params}; +use powersync_sqlite_nostd::ResultCode; +use powersync_sqlite_nostd::bindings::{ + SQLITE_OPEN_CREATE, SQLITE_OPEN_READONLY, SQLITE_OPEN_READWRITE, +}; use serde::Deserialize; +use crate::db::connection::{RawSqliteConnection, SqliteConnection}; use crate::{db::watch::TableNotifiers, error::PowerSyncError}; /// A raw connection pool, giving out both synchronous and asynchronous leases to SQLite @@ -21,46 +21,49 @@ pub struct ConnectionPool { } impl ConnectionPool { - fn prepare_writer(connection: Connection) -> Arc> { + fn prepare_writer(connection: SqliteConnection) -> Arc> { connection - .prepare("SELECT powersync_update_hooks('install');") - .expect("should prepare statement for update hooks") - .query_row(params![], |_| Ok(())) + .exec(c"SELECT powersync_update_hooks('install');") .expect("could not install update hook"); Arc::new(Mutex::new(connection)) } pub fn open>(path: P) -> Result { - let writer = Connection::open(&path)?; + let writer = SqliteConnection::from(RawSqliteConnection::open_path( + &path, + SQLITE_OPEN_READWRITE | SQLITE_OPEN_CREATE, + )?); - writer.pragma_update(None, "journal_mode", "WAL")?; - writer.pragma_update(None, "journal_size_limit", 6 * 1024 * 1024)?; - writer.pragma_update(None, "busy_timeout", 30_000)?; - writer.pragma_update(None, "cache_size", 50 * 1024)?; + writer.exec(c"PRAGMA journal_mode = WAL")?; + writer.exec(c"PRAGMA journal_size_limit = 6291456")?; // 6 * 1024 * 1024 + writer.exec(c"PRAGMA busy_timeout = 30000")?; + writer.exec(c"PRAGMA cache_size = -51200")?; // -(50 * 1024) let mut readers = vec![]; for _ in 0..5 { - let reader = Connection::open(&path)?; - reader.pragma_update(None, "query_only", true)?; - reader.pragma_update(None, "busy_timeout", 30_000)?; + let reader = SqliteConnection::from(RawSqliteConnection::open_path( + &path, + SQLITE_OPEN_READONLY, + )?); + reader.exec(c"PRAGMA query_only = 1")?; readers.push(reader); } Ok(Self::wrap_connections(writer, readers)) } - /// Creates a pool backed by a single write ad multiple reader connections. + /// Creates a pool backed by a single write and multiple reader connections. /// /// Connections will be configured to use WAL mode. pub fn wrap_connections( - writer: Connection, - readers: impl IntoIterator, + writer: impl Into, + readers: impl IntoIterator>, ) -> Self { - let writer = Self::prepare_writer(writer); - let (release, consume) = async_channel::unbounded::(); + let writer = Self::prepare_writer(writer.into()); + let (release, consume) = async_channel::unbounded::(); for reader in readers { - release.send_blocking(reader).unwrap(); + release.send_blocking(reader.into()).unwrap(); } Self { @@ -76,10 +79,10 @@ impl ConnectionPool { } /// Creates a connection pool backed by a single sqlite connection. - pub fn single_connection(conn: Connection) -> Self { + pub fn single_connection(conn: impl Into) -> Self { Self { state: Arc::new(PoolState { - writer: Self::prepare_writer(conn), + writer: Self::prepare_writer(conn.into()), readers: None, table_notifiers: Default::default(), }), @@ -100,7 +103,7 @@ impl ConnectionPool { LeasedConnection { inner: OwnedConnectionLease::Reader { connection: MaybeUninit::new(reader), - release: readers.release_reader.clone(), + pool: self.clone(), }, } } else { @@ -116,19 +119,23 @@ impl ConnectionPool { fn take_update_notifications( &self, - writer: &Connection, - ) -> Result { - let mut stmt = writer.prepare_cached("SELECT powersync_update_hooks('get');")?; - let rows: String = stmt.query_row(params![], |row| row.get(0))?; + writer: &SqliteConnection, + ) -> Result { + let stmt = writer.prepare("SELECT powersync_update_hooks('get');")?; - let updates = serde_json::from_str::(&rows) - .map_err(|_| Error::InvalidQuery)?; + match stmt.step()? { + ResultCode::ROW => { + let updates = + serde_json::from_str::(stmt.column_text(0)?)?; - if !updates.tables.is_empty() { - self.state.table_notifiers.notify_updates(&updates.tables); - } + if !updates.tables.is_empty() { + self.state.table_notifiers.notify_updates(&updates.tables); + } - Ok(updates) + Ok(updates) + } + code => Err(code.into()), + } } async fn take_connection_async(&self, writer: bool) -> LeasedConnection { @@ -142,7 +149,7 @@ impl ConnectionPool { LeasedConnection { inner: OwnedConnectionLease::Reader { connection: MaybeUninit::new(reader), - release: readers.release_reader.clone(), + pool: self.clone(), }, } } else { @@ -180,26 +187,24 @@ pub struct SqliteUpdateNotification { } struct PoolState { - writer: Arc>, + writer: Arc>, readers: Option, table_notifiers: Arc, } struct PoolReaders { - take_reader: Receiver, - release_reader: Sender, + take_reader: Receiver, + release_reader: Sender, } enum OwnedConnectionLease { Writer { - connection: MutexGuardArc, + connection: MutexGuardArc, pool: ConnectionPool, }, Reader { - connection: MaybeUninit, - /// Sender cloned at lease creation — avoids navigating back through the pool - /// and eliminates the need for an `unwrap()` on `pool.state.readers` in Drop. - release: Sender, + connection: MaybeUninit, + pool: ConnectionPool, }, } @@ -210,17 +215,18 @@ impl Drop for OwnedConnectionLease { // Send update notifications for writes made on this connection while leased. let _ = pool.take_update_notifications(connection); } - OwnedConnectionLease::Reader { - connection, - release, - } => { + OwnedConnectionLease::Reader { connection, pool } => { let connection = std::mem::replace(connection, MaybeUninit::uninit()); let connection = unsafe { // safety: Only dropped here connection.assume_init() }; - release + pool.state + .readers + .as_ref() + .unwrap() + .release_reader .send_blocking(connection) .expect("should send connection into pool"); } @@ -235,10 +241,8 @@ pub struct LeasedConnection { inner: OwnedConnectionLease, } -impl Deref for LeasedConnection { - type Target = Connection; - - fn deref(&self) -> &Self::Target { +impl LeasedConnection { + pub(crate) fn sqlite_connection(&self) -> &SqliteConnection { match &self.inner { OwnedConnectionLease::Writer { connection, .. } => connection, OwnedConnectionLease::Reader { connection, .. } => unsafe { @@ -247,10 +251,8 @@ impl Deref for LeasedConnection { }, } } -} -impl DerefMut for LeasedConnection { - fn deref_mut(&mut self) -> &mut Self::Target { + pub(crate) fn sqlite_connection_mut(&mut self) -> &mut SqliteConnection { match &mut self.inner { OwnedConnectionLease::Writer { connection, .. } => connection, OwnedConnectionLease::Reader { connection, .. } => unsafe { @@ -260,3 +262,19 @@ impl DerefMut for LeasedConnection { } } } + +#[cfg(feature = "rusqlite")] +impl Deref for LeasedConnection { + type Target = rusqlite::Connection; + + fn deref(&self) -> &Self::Target { + self.sqlite_connection().rusqlite_connection() + } +} + +#[cfg(feature = "rusqlite")] +impl DerefMut for LeasedConnection { + fn deref_mut(&mut self) -> &mut Self::Target { + self.sqlite_connection_mut().rusqlite_connection_mut() + } +} diff --git a/powersync/src/db/streams.rs b/powersync/src/db/streams.rs index fa6b30a..c31dc78 100644 --- a/powersync/src/db/streams.rs +++ b/powersync/src/db/streams.rs @@ -1,3 +1,4 @@ +use crate::db::connection::{TransactionGuard, exec_stmt}; use crate::{ PowerSyncDatabase, StreamPriority, db::internal::InnerPowerSyncState, @@ -8,7 +9,7 @@ use crate::{ }, util::SerializedJsonObject, }; -use rusqlite::{TransactionBehavior, params}; +use powersync_sqlite_nostd::Destructor; use std::{ cell::Cell, collections::HashMap, @@ -85,14 +86,14 @@ impl<'a> SyncStream<'a> { let serialized = serde_json::to_string(cmd)?; let mut writer = self.db.writer().await?; - let writer = writer.transaction_with_behavior(TransactionBehavior::Immediate)?; + let writer = TransactionGuard::new(writer.sqlite_connection_mut())?; { - let mut stmt = writer.prepare_cached("SELECT powersync_control(?, ?)")?; - let mut rows = stmt.query(params!["subscriptions", serialized])?; - - // Ignore results. - while rows.next()?.is_some() {} + let stmt = writer.inner.prepare("SELECT powersync_control(?, ?)")?; + stmt.bind_text(1, "subscriptions", Destructor::STATIC)?; + // Fine because we drop the statement before serialized + stmt.bind_text(2, &serialized, Destructor::STATIC)?; + exec_stmt(stmt)?; } writer.commit()?; diff --git a/powersync/src/env.rs b/powersync/src/env.rs index daa8e74..edc5d60 100644 --- a/powersync/src/env.rs +++ b/powersync/src/env.rs @@ -1,10 +1,10 @@ -use std::{pin::Pin, time::Duration}; - -use powersync_core::powersync_init_static; - use super::db::pool::ConnectionPool; -use crate::error::PowerSyncError; +use crate::error::{PowerSyncError, RawPowerSyncError}; use crate::http::HttpClient; +use num_traits::FromPrimitive; +use powersync_core::powersync_init_static; +use powersync_sqlite_nostd::ResultCode; +use std::{pin::Pin, time::Duration}; /// All external dependencies required for the PowerSync SDK. /// @@ -38,13 +38,12 @@ impl PowerSyncEnvironment { /// This needs to be invoked before using the PowerSync SDK. It can safely be called multiple /// times. pub fn powersync_auto_extension() -> Result<(), PowerSyncError> { - let rc = powersync_init_static(); - match rc { + match powersync_init_static() { 0 => Ok(()), - _ => Err(rusqlite::Error::SqliteFailure( - rusqlite::ffi::Error::new(rc), - Some("Loading PowerSync core extension failed".into()), - ) + code => Err(RawPowerSyncError::RawSqlite { + code: ResultCode::from_i32(code).unwrap(), + context: "Loading PowerSync core extension failed".into(), + } .into()), } } diff --git a/powersync/src/error.rs b/powersync/src/error.rs index 64a3819..f15f817 100644 --- a/powersync/src/error.rs +++ b/powersync/src/error.rs @@ -1,10 +1,8 @@ +use powersync_sqlite_nostd::ResultCode; use std::error::Error; use std::io; use std::sync::Arc; use std::{borrow::Cow, fmt::Display}; - -use rusqlite::Error as SqliteError; -use rusqlite::types::FromSqlError; use thiserror::Error; pub type Result = std::result::Result; @@ -19,13 +17,23 @@ pub struct PowerSyncError { } impl PowerSyncError { + /// Wrap any error as a PowerSync error to indicate an error in a + /// [crate::BackendConnector::upload_data] implementation. + pub fn upload_error(inner: impl Error + Send + Sync + 'static) -> Self { + RawPowerSyncError::UploadError { + source: Box::new(inner), + } + .into() + } + pub(crate) fn argument_error(desc: impl Into>) -> Self { RawPowerSyncError::ArgumentError { desc: desc.into() }.into() } } -impl From for PowerSyncError { - fn from(value: SqliteError) -> Self { +#[cfg(feature = "rusqlite")] +impl From for PowerSyncError { + fn from(value: rusqlite::Error) -> Self { RawPowerSyncError::Sqlite { inner: value }.into() } } @@ -51,12 +59,6 @@ impl From for PowerSyncError { } } -impl From for PowerSyncError { - fn from(value: io::Error) -> Self { - RawPowerSyncError::IO { inner: value }.into() - } -} - impl Display for PowerSyncError { fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result { self.inner.fmt(f) @@ -74,13 +76,17 @@ pub(crate) enum RawPowerSyncError { #[error("invalid argument: {desc}")] ArgumentError { desc: Cow<'static, str> }, /// An inner SQLite call failed. + #[cfg(feature = "rusqlite")] #[error("SQLite: {inner}")] - Sqlite { inner: SqliteError }, + Sqlite { inner: rusqlite::Error }, + #[error("SQLite: {context} failed with {code}")] + RawSqlite { code: ResultCode, context: String }, /// Reading a value from SQLite failed. + #[cfg(feature = "rusqlite")] #[error("Reading from SQLite: {inner}")] FromSql { #[from] - inner: FromSqlError, + inner: rusqlite::types::FromSqlError, }, /// The version of the core extension linked into the application is unexpected. /// @@ -111,4 +117,19 @@ pub(crate) enum RawPowerSyncError { InvalidCredentials, #[error("Unexpected HTTP status code from PowerSync service: {code}")] UnexpectedStatusCode { code: u16 }, + #[error("Error in upload_data: {source}")] + UploadError { + #[source] + source: Box, + }, +} + +impl From for PowerSyncError { + fn from(value: ResultCode) -> Self { + RawPowerSyncError::RawSqlite { + code: value, + context: String::new(), + } + .into() + } } diff --git a/powersync/src/sync/coordinator.rs b/powersync/src/sync/coordinator.rs index f5b5f12..5a4b47e 100644 --- a/powersync/src/sync/coordinator.rs +++ b/powersync/src/sync/coordinator.rs @@ -88,7 +88,7 @@ impl SyncCoordinator { /// Handle the set of active sync stream subscriptions changing. /// /// This is a no-op if not connected. - pub async fn handle_subscriptions_changed(&self, update: ChangedSyncSubscriptions) { + pub(crate) async fn handle_subscriptions_changed(&self, update: ChangedSyncSubscriptions) { self.download_actor_request(DownloadActorCommand::SubscriptionsChanged(update)) .await; } @@ -118,7 +118,7 @@ impl SyncCoordinator { slot.clone() } - pub fn receive_download_commands(&self) -> Receiver> { + pub(crate) fn receive_download_commands(&self) -> Receiver> { Self::install_actor_channel(&self.control_downloads) } diff --git a/powersync/src/sync/download/actor.rs b/powersync/src/sync/download/actor.rs index 5a57235..8d167ff 100644 --- a/powersync/src/sync/download/actor.rs +++ b/powersync/src/sync/download/actor.rs @@ -21,7 +21,7 @@ use crate::{ }; /// A command sent from a database to the download actor. -pub enum DownloadActorCommand { +pub(crate) enum DownloadActorCommand { Connect(SyncOptions), Disconnect, ResolveOfflineSyncStatusIfNotConnected, @@ -96,7 +96,7 @@ impl DownloadActor { let writer = self.db.writer().await?; self.db .status - .update(|s| s.resolve_offline_state(&writer))?; + .update(|s| s.resolve_offline_state(writer.sqlite_connection()))?; Ok::<(), PowerSyncError>(()) } diff --git a/powersync/src/sync/download/http.rs b/powersync/src/sync/download/http.rs index 1b91c59..1619885 100644 --- a/powersync/src/sync/download/http.rs +++ b/powersync/src/sync/download/http.rs @@ -47,7 +47,7 @@ pub fn sync_stream( let stream = stream::once_future(response); StreamExt::flat_map(stream, |response| match response { - Err(e) => stream::once(Err(e)).boxed(), + Err(error) => stream::once(Err(error)).boxed(), Ok(response) => stream::once(Ok(DownloadEvent::ConnectionEstablished)) .chain(response_to_lines(Ok(response))) .boxed(), @@ -139,3 +139,92 @@ fn response_to_lines( .boxed() } } + +#[cfg(test)] +mod tests { + use std::{pin::Pin, sync::Arc, time::Duration}; + + use async_trait::async_trait; + use futures_lite::{StreamExt, future}; + use rusqlite::Connection; + + use super::*; + use crate::{ + db::{internal::InnerPowerSyncState, pool::ConnectionPool}, + env::{PowerSyncEnvironment, Timer}, + http::{HttpClient, Request, ResponseBody}, + schema::Schema, + sync::coordinator::SyncCoordinator, + }; + + struct FailingClient; + + #[async_trait] + impl HttpClient for FailingClient { + async fn send(&self, _request: Request) -> Result { + Err(PowerSyncError::argument_error("offline")) + } + } + + struct StatusClient(u16); + + #[async_trait] + impl HttpClient for StatusClient { + async fn send(&self, _request: Request) -> Result { + Ok(Response { + status: self.0, + content_type: Some("application/x-ndjson".to_string()), + body: ResponseBody { + reader: stream::empty().boxed(), + length: Some(0), + }, + }) + } + } + + struct UnusedTimer; + + impl Timer for UnusedTimer { + fn delay_once(&self, _duration: Duration) -> Pin + Send>> { + Box::pin(future::pending()) + } + } + + fn first_event(client: impl HttpClient) -> Result, PowerSyncError> { + PowerSyncEnvironment::powersync_auto_extension().unwrap(); + let pool = ConnectionPool::single_connection(Connection::open_in_memory().unwrap()); + let environment = PowerSyncEnvironment::custom(client, pool, &UnusedTimer); + let coordinator = Arc::new(SyncCoordinator::default()); + let db = Arc::new(InnerPowerSyncState::new( + environment, + Schema::default().into(), + &coordinator, + )); + let credentials = PowerSyncCredentials { + endpoint: "https://rust.unit.test.powersync.com/".to_string(), + token: "token".to_string(), + }; + let mut events = Box::pin(sync_stream(db, credentials, "{}".to_string())); + + future::block_on(events.as_mut().try_next()) + } + + #[test] + fn transport_error_does_not_report_connection_established() { + assert!(first_event(FailingClient).is_err()); + } + + #[test] + fn unsuccessful_status_does_not_report_connection_established() { + assert!(first_event(StatusClient(401)).is_err()); + assert!(first_event(StatusClient(500)).is_err()); + } + + #[test] + fn successful_response_reports_connection_before_stream_events() { + assert!(matches!( + first_event(StatusClient(200)), + Ok(Some(DownloadEvent::ConnectionEstablished)) + )); + } +} diff --git a/powersync/src/sync/download/mod.rs b/powersync/src/sync/download/mod.rs index cd95fb1..496b198 100644 --- a/powersync/src/sync/download/mod.rs +++ b/powersync/src/sync/download/mod.rs @@ -2,4 +2,5 @@ mod actor; pub mod http; mod sync_iteration; -pub use actor::{DownloadActor, DownloadActorCommand}; +pub use actor::DownloadActor; +pub(crate) use actor::DownloadActorCommand; diff --git a/powersync/src/sync/download/sync_iteration.rs b/powersync/src/sync/download/sync_iteration.rs index 940c7d2..2d240f5 100644 --- a/powersync/src/sync/download/sync_iteration.rs +++ b/powersync/src/sync/download/sync_iteration.rs @@ -2,13 +2,11 @@ use std::sync::Arc; use futures_lite::{StreamExt, future, stream::Boxed as BoxedStream}; use log::{debug, info, trace, warn}; -use rusqlite::{ - Connection, ToSql, TransactionBehavior, params, - types::{ToSqlOutput, ValueRef}, -}; +use powersync_sqlite_nostd::{Destructor, ManagedStmt, ResultCode}; use serde::Serialize; use serde_json::value::RawValue; +use crate::db::connection::{SqliteConnection, TransactionGuard}; use crate::schema::SchemaOrCustom; use crate::{ SyncOptions, @@ -55,7 +53,7 @@ impl DownloadClient { trace!("Handling event {event:?}"); let instructions = { let mut conn = self.db.writer().await?; - event.invoke_control(&mut conn)? + event.invoke_control(conn.sqlite_connection_mut())? }; for instr in instructions { @@ -184,23 +182,28 @@ impl DownloadEvent { /// Forwards the event to the core extension, and returns instructions that the SDK needs to /// perform. - pub fn invoke_control(self, conn: &mut Connection) -> Result, PowerSyncError> { - let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?; + pub fn invoke_control( + self, + conn: &mut SqliteConnection, + ) -> Result, PowerSyncError> { + let tx = TransactionGuard::new(conn)?; let instructions = { - let mut stmt = tx.prepare_cached("SELECT powersync_control(?, ?)")?; + let stmt = tx.inner.prepare("SELECT powersync_control(?, ?)")?; let (op, arg) = self.into_powersync_control_argument(); - let mut rows = stmt.query(params![op, arg])?; - let Some(row) = rows.next()? else { - return Err(rusqlite::Error::QueryReturnedNoRows)?; - }; + stmt.bind_text(1, op, Destructor::STATIC)?; + arg.bind_to(&stmt, 2)?; - let instructions = row.get_ref(0)?.as_str().map_err(|_| { - PowerSyncError::argument_error("Could not read powersync_control instructions") - })?; + if let ResultCode::ROW = stmt.step()? { + let instructions = stmt.column_text(0).map_err(|_| { + PowerSyncError::argument_error("Could not read powersync_control instructions") + })?; - serde_json::from_str(instructions)? + serde_json::from_str(instructions)? + } else { + panic!("Expected a row") // Can't happen, scalar select + } }; tx.commit()?; @@ -215,14 +218,46 @@ enum PowerSyncControlArgument { Bytes(Vec), } -impl ToSql for PowerSyncControlArgument { - fn to_sql(&self) -> rusqlite::Result> { - Ok(ToSqlOutput::Borrowed(match self { - PowerSyncControlArgument::Null => ValueRef::Null, - PowerSyncControlArgument::StaticString(str) => ValueRef::Text(str.as_bytes()), - PowerSyncControlArgument::String(str) => ValueRef::Text(str.as_bytes()), - PowerSyncControlArgument::Bytes(items) => ValueRef::Blob(items), - })) +impl PowerSyncControlArgument { + fn bind_to(&self, stmt: &ManagedStmt, index: i32) -> Result<(), ResultCode> { + match self { + PowerSyncControlArgument::Null => stmt.bind_null(index), + PowerSyncControlArgument::StaticString(str) => { + stmt.bind_text(index, str, Destructor::STATIC) + } + PowerSyncControlArgument::String(str) => { + stmt.bind_text(index, str, Destructor::TRANSIENT) + } + PowerSyncControlArgument::Bytes(bytes) => { + stmt.bind_blob(index, bytes, Destructor::TRANSIENT) + } + }?; + Ok(()) + } +} + +#[cfg(all(test, feature = "rusqlite"))] +mod tests { + use std::hint::black_box; + + use super::*; + + #[test] + fn binding_dynamic_control_arguments_copies_the_payload() { + let connection = rusqlite::Connection::open_in_memory().unwrap(); + let connection = SqliteConnection::from(connection); + let stmt = connection.prepare("SELECT ?").unwrap(); + let expected = "a".repeat(4096); + + PowerSyncControlArgument::String(expected.clone()) + .bind_to(&stmt, 1) + .unwrap(); + + let overwrite = "b".repeat(4096); + black_box(&overwrite); + + assert_eq!(stmt.step().unwrap(), ResultCode::ROW); + assert_eq!(stmt.column_text(0).unwrap(), expected); } } diff --git a/powersync/src/sync/instruction.rs b/powersync/src/sync/instruction.rs index 9c6785d..6ae570b 100644 --- a/powersync/src/sync/instruction.rs +++ b/powersync/src/sync/instruction.rs @@ -89,7 +89,7 @@ pub struct Timestamp(pub i64); impl From for SystemTime { fn from(val: Timestamp) -> Self { - let since_epoch = Duration::from_secs(val.0 as u64); + let since_epoch = Duration::from_micros(val.0 as u64); SystemTime::UNIX_EPOCH + since_epoch } } diff --git a/powersync/src/sync/status.rs b/powersync/src/sync/status.rs index 16b2a43..20eecb0 100644 --- a/powersync/src/sync/status.rs +++ b/powersync/src/sync/status.rs @@ -7,8 +7,9 @@ use std::{ }; use event_listener::{Event, EventListener}; -use rusqlite::{Connection, params}; +use powersync_sqlite_nostd::ResultCode; +use crate::db::connection::SqliteConnection; use crate::{ error::PowerSyncError, sync::{ @@ -170,12 +171,15 @@ impl SyncStatusData { pub(crate) fn resolve_offline_state( &mut self, - conn: &Connection, + conn: &SqliteConnection, ) -> Result<(), PowerSyncError> { - let mut stmt = conn.prepare_cached("SELECT powersync_offline_sync_status()")?; - let raw_status: String = stmt.query_row(params![], |row| row.get(0))?; + let stmt = conn.prepare("SELECT powersync_offline_sync_status()")?; + let ResultCode::ROW = stmt.step()? else { + panic!("Expected row"); // Can't happen, scalar select + }; - self.update_from_core(serde_json::from_str(&raw_status)?); + let raw_status = stmt.column_text(0)?; + self.update_from_core(serde_json::from_str(raw_status)?); Ok(()) } diff --git a/powersync/src/sync/streams.rs b/powersync/src/sync/streams.rs index a9ec198..8c04f5f 100644 --- a/powersync/src/sync/streams.rs +++ b/powersync/src/sync/streams.rs @@ -124,4 +124,4 @@ impl<'a> From<&'a StreamSubscriptionDescription<'a>> for StreamDescription<'a> { val.description() } } -pub struct ChangedSyncSubscriptions(#[allow(private_interfaces)] pub Vec); +pub(crate) struct ChangedSyncSubscriptions(pub(crate) Vec); diff --git a/powersync/src/sync/upload.rs b/powersync/src/sync/upload.rs index 111f097..955ea39 100644 --- a/powersync/src/sync/upload.rs +++ b/powersync/src/sync/upload.rs @@ -1,12 +1,13 @@ -use std::{collections::HashSet, sync::Arc}; +use std::{collections::HashSet, ops::ControlFlow, sync::Arc}; use futures_lite::{ FutureExt, StreamExt, future::{self, Boxed}, }; use log::{debug, info, warn}; -use rusqlite::{Connection, TransactionBehavior, params}; +use powersync_sqlite_nostd::{Destructor, ResultCode}; +use crate::db::connection::{SqliteConnection, TransactionGuard}; use crate::db::watch::ListenerConfiguration; use crate::sync::coordinator::SyncCoordinator; use crate::{ @@ -59,7 +60,7 @@ impl UploadActor { .env .pool .update_notifiers() - .listen(ListenerConfiguration::if_matches(tables, false)); + .listen(ListenerConfiguration::if_matches(tables, true)); ConnectedUploadActor { connector, crud_stream: stream.map(|_| ()).boxed(), @@ -158,55 +159,16 @@ impl UploadActor { Self::state_transition_from_command_while_uploading(&self.commands, &self.db); let upload_done = async { - let (result, state) = result.await; - - match result { - Ok(_) => { - // It's possible that pending CRUD uploads were preventing data from - // syncing. So now that that's completed, notify the download client in - // case it needs to retry. - if let Some(sync) = self.db.sync.upgrade() { - sync.mark_crud_uploads_completed().await; - } - - // Apart from that, the upload is done and we transition back into the - // ready connected state to start the next iteration when needed. - Some(UploadActorState::Connected(state)) - } - Err(e) => { - warn!("CRUD uploads failed, will retry, {e}"); - self.db - .status - .update(|s| s.set_upload_state(UploadStatus::Error(e))); - let db = self.db.clone(); - - Some(UploadActorState::WaitingForReconnect { - timeout: async move { - db.sync_iteration_delay().await; - state - } - .boxed(), - }) - } - } - }; - - future::race(request, upload_done) - .await - .unwrap_or(old_state) - } - UploadActorState::WaitingForReconnect { ref mut timeout } => { - // Either the timeout expires, in which case we reconnect, or a disconnect is - // requested. - let request = - Self::state_transition_from_command_while_uploading(&self.commands, &self.db); + let state = result.await; + self.db + .status + .update(|s| s.set_upload_state(UploadStatus::Idle)); - let timeout_expired = async { - let state = timeout.await; + // The upload is done and we transition back into the ready connected state to start the next iteration when needed. Some(UploadActorState::Connected(state)) }; - future::race(request, timeout_expired) + future::race(request, upload_done) .await .unwrap_or(old_state) } @@ -222,9 +184,9 @@ impl UploadActor { connector: state.connector.as_ref(), db, }; - let result = upload.run().await; + upload.run().await; - (result, state) + state } .boxed(), } @@ -234,12 +196,7 @@ impl UploadActor { enum UploadActorState { Idle, Connected(ConnectedUploadActor), - RunningUpload { - result: Boxed<(Result<(), PowerSyncError>, ConnectedUploadActor)>, - }, - WaitingForReconnect { - timeout: Boxed, - }, + RunningUpload { result: Boxed }, Stopped, } @@ -262,85 +219,122 @@ struct CrudUpload<'a> { } impl<'a> CrudUpload<'a> { - pub async fn run(&mut self) -> Result<(), PowerSyncError> { + pub async fn run(&mut self) { let mut last_item_id = None::; - while let Some(item) = self.oldest_crud_item_id().await? { - if last_item_id == Some(item) { - warn!("{}", Self::DUPLICATE_ITEM_WARNING); - return Err(PowerSyncError::argument_error( - "Delaying due to previously encountered CRUD item.", - )); + // Invoke upload method on connector until there are no remaining CRUD items to upload. + loop { + match self.upload_step(&mut last_item_id).await { + Ok(ControlFlow::Break(_)) => break, + Ok(ControlFlow::Continue(_)) => continue, + Err(e) => { + last_item_id = None; + info!("CRUD uploads failed, will retry, {e}"); + + self.db + .status + .update(|data| data.set_upload_state(UploadStatus::Error(e))); + self.db.sync_iteration_delay().await; + } } - - last_item_id = Some(item); - self.db - .status - .update(|data| data.set_upload_state(UploadStatus::Uploading)); - self.connector.upload_data().await?; } + } - // Uploading is completed, advance write checkpoint. - if let Some(advance_target) = self.sequence_for_checkpoint().await? { - let write_checkpoint = self.get_write_checkpoint().await?; - advance_target.complete(write_checkpoint, &self.db).await?; + async fn upload_step( + &mut self, + last_item_id: &mut Option, + ) -> Result, PowerSyncError> { + let Some(item) = self.oldest_crud_item_id().await? else { + // Uploading is completed, advance write checkpoint. + if let Some(advance_target) = self.sequence_for_checkpoint().await? { + let write_checkpoint = self.get_write_checkpoint().await?; + advance_target.complete(write_checkpoint, &self.db).await?; + } + + // It's possible that pending CRUD uploads were preventing data from syncing. So now + // that that's completed, notify the download client in case it needs to retry. + if let Some(sync) = self.db.sync.upgrade() { + sync.mark_crud_uploads_completed().await; + } + + return Ok(ControlFlow::Break(())); + }; + + self.db + .status + .update(|data| data.set_upload_state(UploadStatus::Uploading)); + if matches!(*last_item_id, Some(x) if x == item) { + warn!("{}", Self::DUPLICATE_ITEM_WARNING); + return Err(PowerSyncError::argument_error( + "Delaying due to previously encountered CRUD item.", + )); } - Ok(()) + *last_item_id = Some(item); + self.connector.upload_data().await?; + + Ok(ControlFlow::Continue(())) } async fn oldest_crud_item_id(&self) -> Result, PowerSyncError> { let reader = self.db.reader().await?; - Self::read_oldest_crud_item_id(&reader) + Self::read_oldest_crud_item_id(reader.sqlite_connection()) } async fn get_write_checkpoint(&self) -> Result { let client_id = { let reader = self.db.reader().await?; - let mut stmt = reader.prepare("SELECT powersync_client_id()")?; - let id: String = stmt.query_row(params![], |row| row.get(0))?; - id + + let stmt = reader + .sqlite_connection() + .prepare("SELECT powersync_client_id()")?; + let ResultCode::ROW = stmt.step()? else { + panic!("Expected row"); // Can't happen, scalar select + }; + + stmt.column_text(0)?.to_string() }; let credentials = self.connector.fetch_credentials().await?; write_checkpoint(&self.db, &client_id, credentials).await } - fn read_oldest_crud_item_id(conn: &Connection) -> Result, PowerSyncError> { - let mut stmt = conn.prepare("SELECT id FROM ps_crud ORDER BY id LIMIT 1")?; - let mut rows = stmt.query(params![])?; + fn read_oldest_crud_item_id(conn: &SqliteConnection) -> Result, PowerSyncError> { + let stmt = conn.prepare("SELECT id FROM ps_crud ORDER BY id LIMIT 1")?; - Ok(match rows.next()? { - None => None, - Some(row) => Some(row.get(0)?), + Ok(match stmt.step()? { + ResultCode::ROW => Some(stmt.column_int64(0)), + _ => None, }) } - fn ps_crud_sequence(conn: &Connection) -> Result, PowerSyncError> { - let mut seq_before = conn.prepare("SELECT seq FROM main.sqlite_sequence WHERE name = ?")?; - let mut seq_before = seq_before.query(params!["ps_crud"])?; - let Some(row) = seq_before.next()? else { + fn ps_crud_sequence(tx: &TransactionGuard) -> Result, 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 { return Ok(None); }; - Ok(row.get(0)?) + Ok(Some(seq_before.column_int64(0))) } async fn sequence_for_checkpoint( &self, ) -> Result, PowerSyncError> { - let reader = self.db.reader().await?; - { - let mut stmt = - reader.prepare("SELECT 1 FROM ps_buckets WHERE name = ? AND target_op = ?")?; - let mut rows = stmt.query(params!["$local", MAX_OP_ID])?; - if rows.next()?.is_none() { - // Nothing to update. - return Ok(None); - } + let mut reader = self.db.reader().await?; + let reader = reader.sqlite_connection_mut(); + let read_tx = TransactionGuard::new(reader)?; + + 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, })) @@ -366,9 +360,9 @@ impl PendingCheckpointRequest { info!("Updating target to checkpoint {}", self.crud_sequence); let mut writer = db.writer().await?; - let writer = writer.transaction_with_behavior(TransactionBehavior::Immediate)?; + let writer = TransactionGuard::new(writer.sqlite_connection_mut())?; - if CrudUpload::read_oldest_crud_item_id(&writer)?.is_some() { + if CrudUpload::read_oldest_crud_item_id(writer.inner)?.is_some() { warn!("ps_crud is not empty, won't advance target"); return Ok(()); } @@ -384,10 +378,7 @@ impl PendingCheckpointRequest { return Ok(()); } - writer.execute( - "UPDATE ps_buckets SET target_op = ? WHERE name = ?", - params![op_id, "$local"], - )?; + InnerPowerSyncState::target_checkpoint_request_id(&writer, Some(op_id))?; writer.commit()?; Ok(()) diff --git a/powersync/src/util/line_split.rs b/powersync/src/util/line_split.rs index 8452fd7..956a74d 100644 --- a/powersync/src/util/line_split.rs +++ b/powersync/src/util/line_split.rs @@ -41,7 +41,7 @@ impl Stream for LineSplitter { // Split into line including the \n, and the rest let remainder = this.unfinished_line.split_off(idx + 1); let mut completed_line = mem::replace(&mut this.unfinished_line, remainder); - // Remove \n from the completed line, then strip trailing \r for \r\n endings. + // Remove \n, then strip the optional \r from CRLF input. completed_line.pop(); if completed_line.last() == Some(&b'\r') { completed_line.pop(); @@ -95,32 +95,52 @@ mod test { } #[test] - fn utf8_split_across_chunks() { - // "é" is two bytes: 0xC3 0xA9, split across chunk boundary. This verifies we don't try to - // decode individual source chunks. + fn splits_crlf_lines_without_retaining_carriage_returns() { + let bytes = Bytes::copy_from_slice(b"hello\r\nworld\r\n"); + let mut lines = LineSplitter::from(stream::once(Ok(bytes)).boxed()); + + let next = future::block_on(async { lines.try_next().await }).unwrap(); + assert_eq!(next.unwrap(), "hello"); + let next = future::block_on(async { lines.try_next().await }).unwrap(); + assert_eq!(next.unwrap(), "world"); + let next = future::block_on(async { lines.try_next().await }).unwrap(); + assert!(next.is_none()); + } + + #[test] + fn splits_crlf_across_chunks_and_preserves_other_carriage_returns() { let mut lines = LineSplitter::from( stream::iter(vec![ - Ok(Bytes::from_static(b"caf\xc3")), - Ok(Bytes::from_static(b"\xa9\n")), + Ok(Bytes::from_static(b"first\r")), + Ok(Bytes::from_static(b"\nsecond\nlast\r")), ]) .boxed(), ); let next = future::block_on(async { lines.try_next().await }).unwrap(); - assert_eq!(next.unwrap(), "café"); + assert_eq!(next.unwrap(), "first"); + let next = future::block_on(async { lines.try_next().await }).unwrap(); + assert_eq!(next.unwrap(), "second"); + let next = future::block_on(async { lines.try_next().await }).unwrap(); + assert_eq!(next.unwrap(), "last\r"); let next = future::block_on(async { lines.try_next().await }).unwrap(); assert!(next.is_none()); } #[test] - fn splits_crlf_lines() { - let bytes = Bytes::copy_from_slice(b"hello\r\nworld\r\n"); - let mut lines = LineSplitter::from(stream::once(Ok(bytes)).boxed()); + fn utf8_split_across_chunks() { + // "é" is two bytes: 0xC3 0xA9, split across chunk boundary. This verifies we don't try to + // decode individual source chunks. + let mut lines = LineSplitter::from( + stream::iter(vec![ + Ok(Bytes::from_static(b"caf\xc3")), + Ok(Bytes::from_static(b"\xa9\n")), + ]) + .boxed(), + ); let next = future::block_on(async { lines.try_next().await }).unwrap(); - assert_eq!(next.unwrap(), "hello"); - let next = future::block_on(async { lines.try_next().await }).unwrap(); - assert_eq!(next.unwrap(), "world"); + assert_eq!(next.unwrap(), "café"); let next = future::block_on(async { lines.try_next().await }).unwrap(); assert!(next.is_none()); } diff --git a/powersync/tests/crud_test.rs b/powersync/tests/crud_test.rs index 9ad0c3c..be72f6b 100644 --- a/powersync/tests/crud_test.rs +++ b/powersync/tests/crud_test.rs @@ -222,6 +222,35 @@ fn insert() { }); } +#[test] +fn applies_checkpoint_after_draining_crud_queue() { + future::block_on(async move { + let test = DatabaseTest::new(); + let db = test.in_memory_database(); + + execute( + &db, + "INSERT INTO users (id, name) VALUES (?, ?)", + params!["test", "name"], + ) + .await; + + let transaction = db.next_crud_transaction().await.unwrap().unwrap(); + transaction.complete_with_checkpoint(42).await.unwrap(); + + let mut reader = db.reader().await.unwrap(); + let reader = reader.transaction().unwrap(); + let checkpoint: i64 = reader + .query_one( + "SELECT powersync_control('target_checkpoint_request_id', NULL)", + params![], + |row| row.get(0), + ) + .unwrap(); + assert_eq!(checkpoint, 42); + }); +} + #[test] fn crud_transactions() { async fn create_transaction(db: &PowerSyncDatabase, amount: usize) { @@ -303,9 +332,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_row(params![], |_| Ok(())).unwrap(); + stmt.query_one(params![], |_| Ok(())).unwrap(); + drop(stmt); + writer.commit().unwrap(); } assert_eq!( @@ -347,7 +380,7 @@ fn raw_table_crud_trigger() { for write in &["INSERT", "UPDATE", "DELETE"] { trigger_stmt - .query_row( + .query_one( params![serialized_table, format!("users_{write}"), write], |_| Ok(()), ) diff --git a/powersync/tests/database_test.rs b/powersync/tests/database_test.rs index 03770ce..af23bc6 100644 --- a/powersync/tests/database_test.rs +++ b/powersync/tests/database_test.rs @@ -18,7 +18,7 @@ fn link_core_extension() { let reader = db.reader().await.unwrap(); let mut stmt = reader.prepare("SELECT powersync_rs_version();").unwrap(); - let _: String = stmt.query_row(params![], |row| row.get(0)).unwrap(); + let _: String = stmt.query_one(params![], |row| row.get(0)).unwrap(); }); } diff --git a/powersync/tests/sync_test.rs b/powersync/tests/sync_test.rs index dbdc59e..c4e555a 100644 --- a/powersync/tests/sync_test.rs +++ b/powersync/tests/sync_test.rs @@ -1,15 +1,27 @@ +use std::{ + sync::{ + Arc, + atomic::{AtomicUsize, Ordering}, + }, + time::{Duration, SystemTime}, +}; + use async_task::Task; +use async_trait::async_trait; +use event_listener::Event; use futures_lite::{StreamExt, future}; use powersync::{ - PowerSyncDatabase, StreamPriority, StreamSubscription, StreamSubscriptionOptions, SyncOptions, - SyncStatusData, error::PowerSyncError, + BackendConnector, PowerSyncCredentials, PowerSyncDatabase, StreamPriority, StreamSubscription, + StreamSubscriptionOptions, SyncOptions, SyncStatusData, error::PowerSyncError, }; use powersync_test_utils::{ DatabaseTest, mock_sync_service::TestConnector, sync_line::{Checkpoint, SyncLine}, }; +use rusqlite::params; use serde_json::json; +use thiserror::Error; struct SyncStreamTest { test: DatabaseTest, @@ -333,3 +345,208 @@ fn progress_without_priorities() { sync.wait_for_status(|s| !s.is_downloading()).await; }); } + +#[test] +fn upload_retry() { + struct FailOnFirstUpload { + db: PowerSyncDatabase, + counter: Arc, + completed_second: Arc, + } + + #[derive(Error, Debug)] + #[error("Deliberate failure on first upload")] + struct FirstUploadFailure; + + #[async_trait] + impl BackendConnector for FailOnFirstUpload { + async fn fetch_credentials(&self) -> Result { + Ok(PowerSyncCredentials { + endpoint: "https://rust.unit.test.powersync.com/".to_string(), + token: "token".to_string(), + }) + } + + async fn upload_data(&self) -> Result<(), PowerSyncError> { + let Some(tx) = self.db.next_crud_transaction().await? else { + return Ok(()); + }; + + let old_count = self.counter.fetch_add(1, Ordering::SeqCst); + if old_count == 0 { + return Err(PowerSyncError::upload_error(FirstUploadFailure)); + } + + tx.complete().await?; + self.completed_second.notify(usize::MAX); + Ok(()) + } + } + + let sync = SyncStreamTest::new(); + let upload_counter = Arc::new(AtomicUsize::default()); + let event = Arc::new(Event::new()); + let mut options = SyncOptions::new(FailOnFirstUpload { + db: sync.db.clone(), + counter: upload_counter.clone(), + completed_second: event.clone(), + }); + options.with_retry_delay(Duration::ZERO); // We can't use timers in tests + sync.run(sync.db.connect(options)); + + sync.run(async { + sync.wait_for_status(|s| s.is_connected()).await; + + // Trigger a crud upload. + { + let writer = sync.db.writer().await.unwrap(); + writer + .execute( + "INSERT INTO users (id, name) VALUES (uuid(), 'local user')", + params![], + ) + .unwrap(); + } + + // Wait for the second upload to finish. + loop { + let listener = event.listen(); + if upload_counter.load(Ordering::SeqCst) == 2 { + break; + }; + + listener.await + } + + sync.wait_for_status(|s| s.upload_error().is_none() && !s.is_uploading()) + .await; + + assert!(sync.db.next_crud_transaction().await.unwrap().is_none()); + }); +} + +#[test] +fn connect_uploads_crud_that_was_already_queued() { + struct CompleteQueuedUpload { + db: PowerSyncDatabase, + counter: Arc, + } + + #[async_trait] + impl BackendConnector for CompleteQueuedUpload { + async fn fetch_credentials(&self) -> Result { + Ok(PowerSyncCredentials { + endpoint: "https://rust.unit.test.powersync.com/".to_string(), + token: "token".to_string(), + }) + } + + async fn upload_data(&self) -> Result<(), PowerSyncError> { + let Some(transaction) = self.db.next_crud_transaction().await? else { + return Ok(()); + }; + self.counter.fetch_add(1, Ordering::SeqCst); + transaction.complete().await + } + } + + let sync = SyncStreamTest::new(); + sync.run(async { + let writer = sync.db.writer().await.unwrap(); + writer + .execute( + "INSERT INTO users (id, name) VALUES (uuid(), 'queued before connect')", + params![], + ) + .unwrap(); + }); + let upload_counter = Arc::new(AtomicUsize::default()); + sync.run(sync.db.connect(SyncOptions::new(CompleteQueuedUpload { + db: sync.db.clone(), + counter: upload_counter.clone(), + }))); + + sync.run(async { + for _ in 0..100 { + if upload_counter.load(Ordering::SeqCst) != 0 { + break; + } + future::yield_now().await; + } + + assert_eq!(upload_counter.load(Ordering::SeqCst), 1); + assert!(sync.db.next_crud_transaction().await.unwrap().is_none()); + }); +} + +#[test] +fn fetching_credentials_does_not_hold_the_download_writer_lease() { + struct WriterUsingConnector { + entered: async_channel::Sender<()>, + release: async_channel::Receiver<()>, + } + + #[async_trait] + impl BackendConnector for WriterUsingConnector { + async fn fetch_credentials(&self) -> Result { + self.entered.send(()).await.unwrap(); + self.release.recv().await.unwrap(); + Ok(PowerSyncCredentials { + endpoint: "https://rust.unit.test.powersync.com/".to_string(), + token: "token".to_string(), + }) + } + + async fn upload_data(&self) -> Result<(), PowerSyncError> { + Ok(()) + } + } + + let sync = SyncStreamTest::new(); + let (entered_tx, entered_rx) = async_channel::bounded(1); + let (release_tx, release_rx) = async_channel::bounded(1); + sync.run(sync.db.connect(SyncOptions::new(WriterUsingConnector { + entered: entered_tx, + release: release_rx, + }))); + + sync.run(async { + entered_rx.recv().await.unwrap(); + let writer = future::poll_once(sync.db.writer()).await; + assert!( + writer.is_some(), + "download retained the writer while awaiting credentials" + ); + drop(writer); + release_tx.send(()).await.unwrap(); + }); +} + +#[test] +fn reports_correct_times() { + let sync = SyncStreamTest::new(); + sync.connect(); + + sync.run(async { + let request = sync.test.http.receive_requests.recv().await.unwrap(); + sync.wait_for_status(|s| s.is_connected()).await; + + request + .send_checkpoint(Checkpoint::single_bucket("a", 0, None)) + .await; + request.send_checkpoint_complete(0, None).await; + sync.wait_for_status(|s| !s.is_downloading()).await; + + let stream = sync.db.sync_stream("a", None); + let status = sync.db.status(); + let status = status + .for_stream(&stream) + .expect("should have stream status"); + let last_synced_at = status + .subscription + .last_synced_at() + .expect("should have last synced at"); + let delta = SystemTime::now().duration_since(last_synced_at).unwrap(); + assert!(delta < Duration::from_secs(5)); + }); +} diff --git a/powersync_test_utils/Cargo.toml b/powersync_test_utils/Cargo.toml index 2591c52..9e801b2 100644 --- a/powersync_test_utils/Cargo.toml +++ b/powersync_test_utils/Cargo.toml @@ -16,7 +16,7 @@ log = "0.4.28" bytes = "1" pin-project-lite = "0.2.16" powersync = { path = "../powersync" } -rusqlite = { version = "0.32.0", features = ["load_extension", "bundled"] } +rusqlite = { version = "0.39.0", features = ["load_extension", "bundled"] } serde = "1.0.228" serde_json = "1.0.145" serde_with = "3.15.0"