Skip to content
Draft
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
6 changes: 6 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,3 +1,9 @@
## 0.2.0 (unreleased)

- __Breaking__: Split `BackendConnector` trait into `Authenticator` and `MutationUploader`.
- `Authenticator` is allowed to cache tokens now, and is informed when a token has expired.
- Add `SyncOptions::download_only` and `SyncOptions::upload_only` for unidirectional sync.

## 0.1.0

- __Breaking__: `PowerSyncEnvironment::custom` now needs a way to spawn async
Expand Down
49 changes: 36 additions & 13 deletions examples/egui_todolist/src/database.rs
Original file line number Diff line number Diff line change
@@ -1,8 +1,10 @@
use std::sync::{Arc, Mutex};

use async_trait::async_trait;
use futures_lite::StreamExt;
use log::warn;
use powersync::{
BackendConnector, CheckpointMode, ConnectionPool, PowerSyncCredentials, PowerSyncDatabase,
Authenticator, CheckpointMode, ConnectionPool, MutationUploader, PowerSyncDatabase,
SyncOptions, UpdateType,
env::PowerSyncEnvironment,
error::PowerSyncError,
Expand Down Expand Up @@ -79,6 +81,7 @@ impl TodoList {
#[derive(Clone)]
pub struct TodoDatabase {
pub db: PowerSyncDatabase,
cached_credentials: Arc<Mutex<Option<Arc<String>>>>,
}

impl TodoDatabase {
Expand All @@ -94,11 +97,14 @@ impl TodoDatabase {
schema.tables.push(TodoEntry::schema());

let db = PowerSyncDatabase::new(env, schema);
Self { db }
Self {
db,
cached_credentials: Arc::default(),
}
}

pub async fn connect(&self) {
let mut options = SyncOptions::new(self.clone());
let mut options = SyncOptions::new("http://localhost:8080", self.clone(), self.clone());
options.with_checkpoint_mode(CheckpointMode::Requests(Default::default()));
self.db.connect(options).await
}
Expand All @@ -107,29 +113,46 @@ impl TodoDatabase {
self.db.disconnect().await;
}

async fn fetch_credentials_self_hosted(&self) -> Result<PowerSyncCredentials, PowerSyncError> {
async fn fetch_credentials_self_hosted(&self) -> Result<Arc<String>, PowerSyncError> {
let response = reqwest::get("http://localhost:6060/api/auth/token").await?;

#[derive(Deserialize)]
struct TokenResponse {
token: String,
token: Arc<String>,
}

let token: TokenResponse = response.json().await?;
Ok(PowerSyncCredentials {
endpoint: "http://localhost:8080".to_string(),
token: token.token,
})
Ok(token.token)
}
}

#[async_trait]
impl BackendConnector for TodoDatabase {
async fn fetch_credentials(&self) -> Result<PowerSyncCredentials, PowerSyncError> {
self.fetch_credentials_self_hosted().await
impl Authenticator for TodoDatabase {
async fn resolve_credentials(&self) -> Result<Arc<String>, PowerSyncError> {
{
let guard = self.cached_credentials.lock().unwrap();
if let Some(ref cached) = *guard {
// Fast path, we have cached credentials. Concurrent fetch_credentials_self_hosted
// are okay too.
return Ok(cached.clone());
}
}

let credentials = self.fetch_credentials_self_hosted().await?;
let mut guard = self.cached_credentials.lock().unwrap();
*guard = Some(credentials.clone());
Ok(credentials)
}

async fn upload_data(&self) -> Result<(), PowerSyncError> {
fn invalidate_credentials(&self) {
let mut guard = self.cached_credentials.lock().unwrap();
*guard = None;
}
}

#[async_trait]
impl MutationUploader for TodoDatabase {
async fn upload(&self) -> Result<(), PowerSyncError> {
let mut transactions = self.db.crud_transactions();
let mut last_tx = None;

Expand Down
10 changes: 9 additions & 1 deletion powersync/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -61,7 +61,15 @@ stream transformer.
### Connecting

To automatically keep the local SQLite database in-sync with a backend process, call `PowerSyncDatabase::connect`.
This requires passing your own `BackendConnector` implementation, see the examples for a possible implementation.
This requires passing `SyncOptions`, which can be created for:

- downloading only (via `SyncOptions::download_only`), which requires the endpoint of your PowerSync instance and an
`powersync::Authenticator` implementation returning JWT tokens for the client to use.
- uploading only (via `SyncOptions::upload_only`), which requires a `powersync::MutationUploader` responsible for
uploading local mutations to your backend. In this mode, no connection to a PowerSync service is made.
- bidirectional sync (recommended, via `SyncOptions::new`). This requires both an authenticator and a mutation uploader.

Also see the examples for a possible implementation.

### Sync status

Expand Down
2 changes: 1 addition & 1 deletion powersync/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@ pub use db::streams::StreamSubscriptionOptions;
pub use db::streams::SyncStream;
pub use db::{DisconnectAndClearFlags, PowerSyncDatabase};
pub use sync::checkpoint::{CheckpointError, CheckpointRequest};
pub use sync::connector::{BackendConnector, PowerSyncCredentials};
pub use sync::connector::{Authenticator, MutationUploader};
pub use sync::options::{CheckpointMode, RequestsCheckpointMode, SyncOptions};
pub use sync::status::SyncStatusData;
pub use sync::stream_priority::StreamPriority;
Expand Down
20 changes: 10 additions & 10 deletions powersync/src/sync/checkpoint.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,13 +5,14 @@ use log::{debug, warn};
use thiserror::Error;

use crate::{
BackendConnector, CheckpointMode, RequestsCheckpointMode, SyncOptions, SyncStatusData,
CheckpointMode, RequestsCheckpointMode, SyncOptions, SyncStatusData,
db::internal::InnerPowerSyncState,
error::{PowerSyncError, RawPowerSyncError},
sync::{
coordinator::{SyncChannels, SyncCoordinator},
download::http::checkpoint_request,
instruction::CheckpointRequestPayload,
options::EndpointAndAuthenticator,
upload::get_client_id,
},
};
Expand Down Expand Up @@ -113,16 +114,15 @@ pub async fn repost_unacknowledged_checkpoints(
let CheckpointMode::Requests(requests) = options.checkpoints else {
return;
};
let Some(ref authenticator) = options.endpoint else {
return;
};

loop {
// Make sure the system is seeded and ready.
let result = repost_unacknowledged_checkpoint_iteration(
&db,
&channels,
options.connector.as_ref(),
requests,
)
.await;
let result =
repost_unacknowledged_checkpoint_iteration(&db, &channels, authenticator, requests)
.await;

if let Err(err) = result {
if let RawPowerSyncError::Checkpoint {
Expand All @@ -141,7 +141,7 @@ pub async fn repost_unacknowledged_checkpoints(
async fn repost_unacknowledged_checkpoint_iteration(
db: &InnerPowerSyncState,
channels: &SyncChannels,
connector: &dyn BackendConnector,
authenticator: &EndpointAndAuthenticator,
mode: RequestsCheckpointMode,
) -> Result<(), PowerSyncError> {
// Make sure the system is seeded and ready
Expand Down Expand Up @@ -187,7 +187,7 @@ async fn repost_unacknowledged_checkpoint_iteration(
debug!("Retrying checkpoint request id {request_id}");
checkpoint_request(
&db,
connector,
authenticator,
&CheckpointRequestPayload {
client_id: get_client_id(db).await?,
checkpoint_request_id: request_id,
Expand Down
60 changes: 45 additions & 15 deletions powersync/src/sync/connector.rs
Original file line number Diff line number Diff line change
@@ -1,19 +1,26 @@
use std::pin::Pin;
use std::{pin::Pin, sync::Arc};

use async_trait::async_trait;
use url::Url;

use crate::error::{PowerSyncError, RawPowerSyncError};

/// A backend connector is responsible for upload local writes as well as resolving JWTs used to
/// connect to the PowerSync service.
/// Authenticates the PowerSync SDK against a PowerSync service, allowing it to sync changes from a
/// backend source database.
#[async_trait]
pub trait BackendConnector: Send + Sync {
/// Fetches a fresh JWT from the backend to be used against the PowerSync service.
async fn fetch_credentials(&self) -> Result<PowerSyncCredentials, PowerSyncError>;
pub trait Authenticator: Send + Sync {
/// Resolves a JWT to use when connecting to a PowerSync service.
///
/// The SDK does not cache this value and will call this method multiple times while a
/// database is connected. Implementations should consider caching credentials.
async fn resolve_credentials(&self) -> Result<Arc<String>, PowerSyncError>;

/// Inspects completed CRUD transactions on a database and uploads them.
async fn upload_data(&self) -> Result<(), PowerSyncError>;
/// Invoked by the SDK when the PowerSync service has rejected credentials previously
/// returned by [Self::resolve_credentials].
///
/// If the ocnnector returns cached tokens, it can use this as a hint to refresh its local
/// state.
fn invalidate_credentials(&self) {}

/// This is optional, and should only return a future for connectors capable of requesting
/// checkpoints.
Expand All @@ -36,15 +43,37 @@ pub trait BackendConnector: Send + Sync {
}
}

/// Uploads local mutations (from `INSERT`, `UPDATE` and `DELETE` statements against the
/// local database) to your source database.
///
/// While simple cases might use a protocol like PostgREST to write directly into the database,
/// a custom backend is commonly used to validate uploaded mutations.
#[async_trait]
pub trait MutationUploader: Send + Sync {
async fn upload(&self) -> Result<(), PowerSyncError>;
}

/// Allows using async functions and closures (e.g. `|| async { ... }`) as a [MutationUploader].
#[async_trait]
impl<F, Fut> MutationUploader for F
where
F: Fn() -> Fut + Send + Sync,
Fut: Future<Output = Result<(), PowerSyncError>> + Send,
{
async fn upload(&self) -> Result<(), PowerSyncError> {
self().await
}
}

/// Credentials used to connect to a PowerSync service instance.
pub struct PowerSyncCredentials {
pub struct PowerSyncCredentials<'a> {
/// PowerSync endpoint, e.g. `https://myinstance.powersync.co`.
pub endpoint: String,
pub endpoint: &'a str,
/// The token used to authenticate against the PowerSync service.
pub token: String,
pub token: Arc<String>,
}

impl PowerSyncCredentials {
impl<'a> PowerSyncCredentials<'a> {
/// Parses the [Self::endpoint] into a URI.
pub(crate) fn parsed_endpoint(&self, endpoint: &str) -> Result<Url, PowerSyncError> {
let url = Url::parse(&self.endpoint)
Expand All @@ -61,12 +90,13 @@ impl PowerSyncCredentials {

#[cfg(test)]
mod test {
use crate::PowerSyncCredentials;
use super::PowerSyncCredentials;
use std::sync::Arc;

fn is_endpoint_valid(endpoint: &str) -> bool {
PowerSyncCredentials {
token: "".to_string(),
endpoint: endpoint.to_string(),
token: Arc::new("test".to_string()),
endpoint,
}
.parsed_endpoint("")
.is_ok()
Expand Down
7 changes: 5 additions & 2 deletions powersync/src/sync/coordinator.rs
Original file line number Diff line number Diff line change
Expand Up @@ -166,9 +166,12 @@ impl SyncCoordinator {
) -> Result<CheckpointRequest, CheckpointError> {
let guard = self.task.lock().await;
let tasks = Self::extract_connected_with_requests(guard.as_ref())?;
let Some(ref authenticator) = tasks.options.endpoint else {
return Err(CheckpointError::Disconnected);
};

let authenticator = authenticator.clone();
let channels = tasks.channels.clone();
let connector = tasks.options.connector.clone();
// Avoid holding the lock across a suspension point. It's fine if there's a concurrent
// reconnect, post_checkpoint_request will return an error in that case.
drop(guard);
Expand All @@ -177,7 +180,7 @@ impl SyncCoordinator {
.await
.map_err(CheckpointError::as_request_error)?;
let checkpoint_request_id =
post_checkpoint_request(client_id, connector.as_ref(), &channels, &db)
post_checkpoint_request(client_id, &authenticator, &channels, &db)
.await
.map_err(CheckpointError::as_request_error)?;

Expand Down
Loading
Loading