use std::collections::VecDeque;
use serde_json::{json, Value};
use syncular_client::{
ClientDiagnosticsRequest, ClientDiagnosticsSnapshot, SyncClient, SyncIntent,
};
use syncular_command::{dispatch, CreateEffects};
use crate::transport::{self, HostTransport};
#[derive(Debug, Clone)]
pub struct Event {
pub json: Value,
}
pub struct SyncularCore {
client: Option<SyncClient>,
transport: HostTransport,
effects: CreateEffects,
queue: VecDeque<Event>,
last_diagnostics_snapshot: Option<ClientDiagnosticsSnapshot>,
diagnostics_observed: bool,
interactive_sync: bool,
background_sync_ms: Option<u64>,
}
impl SyncularCore {
pub fn new(config: &Value) -> Result<Self, String> {
Self::new_with_notify(config, None)
}
pub fn new_with_notify(
config: &Value,
notify: Option<std::sync::Arc<dyn Fn() + Send + Sync>>,
) -> Result<Self, String> {
let transport = HostTransport::from_config_with_notify(config, notify)?;
Ok(SyncularCore {
client: None,
transport,
effects: CreateEffects::default(),
queue: VecDeque::new(),
last_diagnostics_snapshot: None,
diagnostics_observed: false,
interactive_sync: false,
background_sync_ms: None,
})
}
pub fn command(&mut self, command: &Value) -> Value {
let method = command.get("method").and_then(Value::as_str).unwrap_or("");
let params = command.get("params").cloned().unwrap_or(Value::Null);
if method == "enableDiagnostics" {
self.diagnostics_observed = true;
self.last_diagnostics_snapshot = None;
self.drain_realtime();
self.drain_core_outputs();
self.emit_diagnostics_if_changed();
return json!({ "result": {} });
}
if method == "diagnosticsSnapshot" {
self.diagnostics_observed = true;
}
let result = dispatch(
&mut self.transport,
&mut self.client,
&mut self.effects,
method,
¶ms,
);
if method == "create" {
self.last_diagnostics_snapshot = None;
self.transport.set_signed_urls(self.effects.signed_urls);
}
if method == "beginSecurityPreflight"
|| method == "shutdown"
|| (method == "create"
&& params
.get("securityPreflight")
.and_then(Value::as_bool)
.unwrap_or(false))
{
self.interactive_sync = false;
self.background_sync_ms = None;
}
if let Ok(value) = &result {
if value.pointer("/effects/sync/kind").and_then(Value::as_str) == Some("interactive") {
self.interactive_sync = true;
}
}
self.drain_realtime();
self.drain_core_outputs();
self.emit_diagnostics_if_changed();
match result {
Ok(mut value) => {
if let Some(object) = value.as_object_mut() {
object.remove("effects");
}
json!({ "result": value })
}
Err((code, message)) => json!({ "error": { "code": code, "message": message } }),
}
}
pub fn query(&mut self, sql: &str, params: Value) -> Value {
let bind = match params {
Value::Null => Value::Array(Vec::new()),
other => other,
};
self.command(&json!({ "method": "query", "params": { "sql": sql, "params": bind } }))
}
pub fn take_sync_intent(&mut self) -> SyncIntent {
if std::mem::take(&mut self.interactive_sync) {
self.background_sync_ms = None;
SyncIntent::Interactive
} else if let Some(delay_ms) = self.background_sync_ms.take() {
SyncIntent::Background { delay_ms }
} else {
SyncIntent::None
}
}
pub fn sync_until_idle(&mut self) -> Value {
if self.client.is_none() {
return json!({ "result": null });
}
self.command(&json!({ "method": "syncUntilIdle", "params": {} }))
}
pub fn poll_transport(&mut self) {
self.drain_realtime();
self.drain_core_outputs();
self.emit_diagnostics_if_changed();
}
pub fn drain_events(&mut self) -> Vec<Event> {
self.queue.drain(..).collect()
}
pub fn set_headers(&mut self, headers: Vec<(String, String)>) {
self.transport.set_headers(headers);
}
pub fn shutdown(&mut self) {
if let Some(mut client) = self.client.take() {
client.disconnect_realtime(&mut self.transport);
client.seal_security_on_teardown();
}
self.interactive_sync = false;
self.background_sync_ms = None;
self.diagnostics_observed = false;
self.transport.shutdown();
}
fn push(&mut self, json: Value) {
self.queue.push_back(Event { json });
}
fn drain_realtime(&mut self) {
if self.client.is_none() {
return;
}
let frames = self.transport.take_inbound();
for frame in frames {
match frame {
transport::Inbound::Text(text) => {
if is_presence_control(&text) {
self.push(json!({ "type": "presence" }));
}
if let Some(client) = self.client.as_mut() {
client.on_realtime_text(&text);
}
}
transport::Inbound::Binary(bytes) => {
if let Some(client) = self.client.as_mut() {
client.on_realtime_binary(&mut self.transport, &bytes);
}
}
}
}
}
fn drain_core_outputs(&mut self) {
let Some(client) = self.client.as_mut() else {
return;
};
let batches = client.drain_change_batches();
let intents = client.drain_sync_intents();
for batch in batches {
self.push(json!({ "type": "change", "batch": batch }));
}
for intent in intents {
match intent {
SyncIntent::Interactive => self.interactive_sync = true,
SyncIntent::Background { delay_ms } => {
self.background_sync_ms = Some(
self.background_sync_ms
.map_or(delay_ms, |current| current.min(delay_ms)),
);
}
SyncIntent::None => {}
}
}
}
fn emit_diagnostics_if_changed(&mut self) {
if !self.diagnostics_observed {
return;
}
let Some(client) = self.client.as_ref() else {
return;
};
if client.security_preflight() {
return;
}
let Ok(snapshot) = client.diagnostics_snapshot(&ClientDiagnosticsRequest::default()) else {
return;
};
if let Some(previous) = &mut self.last_diagnostics_snapshot {
previous.captured_at_ms = snapshot.captured_at_ms;
if previous == &snapshot {
return;
}
}
let event = json!({ "type": "diagnostics", "snapshot": snapshot });
self.last_diagnostics_snapshot = Some(snapshot);
self.push(event);
}
}
fn is_presence_control(text: &str) -> bool {
serde_json::from_str::<Value>(text)
.ok()
.and_then(|v| {
v.get("event")
.and_then(Value::as_str)
.map(|e| e == "presence")
})
.unwrap_or(false)
}
#[cfg(test)]
mod tests {
use super::*;
fn simple_schema() -> Value {
json!({
"version": 1,
"tables": [{
"name": "todo",
"primaryKey": "id",
"columns": [
{ "name": "id", "type": "string", "nullable": false },
{ "name": "title", "type": "string", "nullable": false },
{ "name": "done", "type": "boolean", "nullable": false }
],
"scopes": []
}]
})
}
fn create(core: &mut SyncularCore) {
let reply = core.command(&json!({
"method": "create",
"params": { "clientId": "c1", "schema": simple_schema() }
}));
assert_eq!(reply["result"], json!({}), "create ok: {reply}");
}
#[test]
fn command_round_trip_create_mutate_query() {
let mut core = SyncularCore::new(&json!({})).unwrap();
create(&mut core);
let sub = core.command(&json!({
"method": "subscribe",
"params": { "id": "s1", "table": "todo", "scopes": {} }
}));
assert_eq!(sub["result"], json!({}));
let mutate = core.command(&json!({
"method": "mutate",
"params": { "mutations": [{
"op": "upsert", "table": "todo",
"values": { "id": "t1", "title": "hello", "done": false }
}] }
}));
assert!(mutate["result"]["clientCommitId"].is_string(), "{mutate}");
let rows = core.query("SELECT id, title FROM todo ORDER BY id", Value::Null);
let list = rows["result"]["rows"].as_array().expect("rows");
assert_eq!(list.len(), 1);
assert_eq!(list[0]["title"], "hello");
assert_eq!(list[0]["id"], "t1");
}
#[test]
fn query_binds_params() {
let mut core = SyncularCore::new(&json!({})).unwrap();
create(&mut core);
core.command(&json!({
"method": "mutate",
"params": { "mutations": [
{ "op": "upsert", "table": "todo", "values": { "id": "a", "title": "A", "done": false } },
{ "op": "upsert", "table": "todo", "values": { "id": "b", "title": "B", "done": true } }
] }
}));
let rows = core.query("SELECT id FROM todo WHERE done = ?", json!([true]));
let list = rows["result"]["rows"].as_array().expect("rows");
assert_eq!(list.len(), 1);
assert_eq!(list[0]["id"], "b");
}
#[test]
fn events_derived_after_mutate() {
let mut core = SyncularCore::new(&json!({})).unwrap();
create(&mut core);
let enabled = core.command(&json!({ "method": "enableDiagnostics", "params": {} }));
assert_eq!(enabled["result"], json!({}));
let _ = core.drain_events();
core.command(&json!({
"method": "mutate",
"params": { "mutations": [{
"op": "upsert", "table": "todo",
"values": { "id": "t1", "title": "x", "done": false }
}] }
}));
let events = core.drain_events();
let kinds: Vec<&str> = events
.iter()
.filter_map(|e| e.json.get("type").and_then(Value::as_str))
.collect();
assert!(kinds.contains(&"change"), "kinds: {kinds:?}");
assert!(kinds.contains(&"diagnostics"), "kinds: {kinds:?}");
let change = events
.iter()
.find(|event| event.json["type"] == "change")
.expect("change event");
assert_eq!(change.json["batch"]["revision"], "1");
assert_eq!(change.json["batch"]["tables"][0]["table"], "todo");
assert_eq!(change.json["batch"]["status"]["outbox"], 1);
assert_eq!(change.json["batch"]["status"]["syncNeeded"], false);
assert!(matches!(core.take_sync_intent(), SyncIntent::Interactive));
assert!(core.drain_events().is_empty());
}
#[test]
fn diagnostics_keep_fresh_capture_times_and_revision_independent_changes() {
let mut core = SyncularCore::new(&json!({})).unwrap();
let created = core.command(&json!({"method": "create", "params": {
"clientId": "diagnostics-comparison", "schema": simple_schema(), "nowMs": 1000
}}));
assert!(created.get("error").is_none(), "{created}");
assert!(core.drain_events().is_empty());
core.command(&json!({"method": "enableDiagnostics", "params": {}}));
let initial = core.drain_events();
assert_eq!(initial.len(), 1);
assert_eq!(initial[0].json["snapshot"]["capturedAtMs"], 1000);
core.client.as_mut().unwrap().set_now_ms(2000);
let read = json!({"method": "query", "params": {"sql": "SELECT 1 AS id"}});
assert_eq!(core.command(&read)["result"]["rows"], json!([{"id": 1}]));
assert!(core.drain_events().is_empty());
let changed = core.command(&json!({"method": "mutate", "params": {
"mutations": [{"op": "upsert", "table": "todo", "values": {
"id": "first", "title": "private value", "done": false
}}]
}}));
assert!(changed.get("error").is_none(), "{changed}");
let expected = serde_json::to_value(
core.client
.as_ref()
.unwrap()
.diagnostics_snapshot(&Default::default())
.unwrap(),
)
.unwrap();
let diagnostics: Vec<Value> = core
.drain_events()
.into_iter()
.filter(|event| event.json["type"] == "diagnostics")
.map(|event| event.json)
.collect();
assert_eq!(
diagnostics,
vec![json!({"type": "diagnostics", "snapshot": expected})]
);
assert_eq!(diagnostics[0]["snapshot"]["capturedAtMs"], 2000);
core.client.as_mut().unwrap().set_now_ms(3000);
assert!(core.command(&read).get("error").is_none());
assert!(core.drain_events().is_empty());
let revision = core.client.as_ref().unwrap().local_revision();
core.command(&json!({"method": "sync", "params": {}}));
assert_eq!(core.client.as_ref().unwrap().local_revision(), revision);
let expected = serde_json::to_value(
core.client
.as_ref()
.unwrap()
.diagnostics_snapshot(&Default::default())
.unwrap(),
)
.unwrap();
let diagnostics: Vec<Value> = core
.drain_events()
.into_iter()
.filter(|event| event.json["type"] == "diagnostics")
.map(|event| event.json)
.collect();
assert_eq!(
diagnostics,
vec![json!({"type": "diagnostics", "snapshot": expected})]
);
assert_eq!(diagnostics[0]["snapshot"]["capturedAtMs"], 3000);
assert_eq!(diagnostics[0]["snapshot"]["lastRound"]["status"], "failed");
let preflight = core.command(&json!({"method": "beginSecurityPreflight", "params": {}}));
assert!(preflight.get("error").is_none(), "{preflight}");
assert_eq!(
core.command(&read)["error"]["code"],
"client.security_preflight_required"
);
assert!(core
.drain_events()
.iter()
.all(|event| event.json["type"] != "diagnostics"));
}
#[test]
fn diagnostics_events_wait_for_a_registered_consumer() {
let mut core = SyncularCore::new(&json!({})).unwrap();
create(&mut core);
let _ = core.drain_events();
core.command(&json!({
"method": "mutate",
"params": { "mutations": [{
"op": "upsert", "table": "todo",
"values": { "id": "t1", "title": "x", "done": false }
}] }
}));
let kinds: Vec<String> = core
.drain_events()
.iter()
.filter_map(|e| e.json.get("type").and_then(Value::as_str))
.map(str::to_owned)
.collect();
assert!(kinds.contains(&"change".to_owned()), "kinds: {kinds:?}");
assert!(
!kinds.contains(&"diagnostics".to_owned()),
"kinds: {kinds:?}"
);
let enabled = core.command(&json!({ "method": "enableDiagnostics", "params": {} }));
assert_eq!(enabled["result"], json!({}));
let events = core.drain_events();
assert!(
events.iter().any(|e| e.json["type"] == "diagnostics"),
"events: {events:?}"
);
core.command(&json!({
"method": "mutate",
"params": { "mutations": [{
"op": "upsert", "table": "todo",
"values": { "id": "t2", "title": "y", "done": false }
}] }
}));
let events = core.drain_events();
assert!(
events.iter().any(|e| e.json["type"] == "diagnostics"),
"events: {events:?}"
);
}
#[test]
fn a_snapshot_pull_registers_the_diagnostics_consumer() {
let mut core = SyncularCore::new(&json!({})).unwrap();
create(&mut core);
let _ = core.drain_events();
let reply = core.command(&json!({ "method": "diagnosticsSnapshot", "params": {} }));
assert_eq!(reply["result"]["version"], 1);
let _ = core.drain_events();
core.command(&json!({
"method": "mutate",
"params": { "mutations": [{
"op": "upsert", "table": "todo",
"values": { "id": "t1", "title": "x", "done": false }
}] }
}));
let events = core.drain_events();
assert!(
events.iter().any(|e| e.json["type"] == "diagnostics"),
"events: {events:?}"
);
}
#[test]
fn diagnostics_are_versioned_bounded_and_payload_free() {
let mut core = SyncularCore::new(&json!({})).unwrap();
create(&mut core);
let reply = core.command(&json!({
"method": "diagnosticsSnapshot",
"params": {
"expectedSubscriptions": [{ "id": "membership", "table": "todo" }]
}
}));
assert_eq!(reply["result"]["version"], 1);
assert_eq!(reply["result"]["subscriptions"][0]["state"], "unregistered");
let encoded = reply.to_string();
assert!(!encoded.contains("clientId"));
assert!(!encoded.contains("dbPath"));
assert!(!encoded.contains("operations"));
}
#[test]
fn sync_without_native_transport_fails_loud() {
let mut core = SyncularCore::new(&json!({})).unwrap();
create(&mut core);
let outcome = core.command(&json!({ "method": "sync", "params": {} }));
assert_eq!(outcome["result"]["ok"], json!(false), "{outcome}");
assert_eq!(outcome["result"]["errorCode"], "transport.unavailable");
}
#[test]
fn file_db_persists_across_reopen() {
let dir = std::env::temp_dir();
let path = dir.join(format!("syncular-tauri-test-{}.db", std::process::id()));
let path_str = path.to_string_lossy().to_string();
let _ = std::fs::remove_file(&path);
{
let mut core = SyncularCore::new(&json!({})).unwrap();
let reply = core.command(&json!({
"method": "create",
"params": { "clientId": "c1", "schema": simple_schema(), "dbPath": path_str }
}));
assert_eq!(reply["result"], json!({}), "create with dbPath: {reply}");
core.command(&json!({
"method": "mutate",
"params": { "mutations": [{
"op": "upsert", "table": "todo",
"values": { "id": "persisted", "title": "kept", "done": false }
}] }
}));
let revision = core.command(&json!({
"method": "localRevision", "params": {}
}));
assert_eq!(revision["result"]["revision"], "1");
}
{
let mut core = SyncularCore::new(&json!({})).unwrap();
let reopened = core.command(&json!({
"method": "create",
"params": { "schema": simple_schema(), "dbPath": path_str }
}));
assert_eq!(reopened["result"], json!({}), "reopen: {reopened}");
let rows = core.query("SELECT title FROM todo", Value::Null);
let list = rows["result"]["rows"].as_array().expect("rows");
assert_eq!(list.len(), 1, "reopened db: {rows}");
assert_eq!(list[0]["title"], "kept");
let revision = core.command(&json!({
"method": "localRevision", "params": {}
}));
assert_eq!(revision["result"]["revision"], "1");
let pending = core.command(&json!({
"method": "pendingCommitIds", "params": {}
}));
assert_eq!(pending["result"]["ids"].as_array().map(Vec::len), Some(1));
let status = core.command(&json!({
"method": "statusSnapshot", "params": {}
}));
assert_eq!(status["result"]["outbox"], 1);
assert_eq!(status["result"]["syncNeeded"], true);
assert!(matches!(core.take_sync_intent(), SyncIntent::Interactive));
}
{
let mut core = SyncularCore::new(&json!({})).unwrap();
let mismatch = core.command(&json!({
"method": "create",
"params": { "clientId": "different", "schema": simple_schema(), "dbPath": path_str }
}));
assert_eq!(mismatch["error"]["code"], "client.identity_mismatch");
}
let _ = std::fs::remove_file(&path);
}
#[test]
fn config_validation_rejects_baseurl_without_native_feature() {
let result = SyncularCore::new(&json!({ "baseUrl": "http://localhost:9/sync" }));
#[cfg(not(feature = "native-transport"))]
assert!(
result.is_err(),
"baseUrl must be refused without native-transport"
);
#[cfg(feature = "native-transport")]
assert!(result.is_ok(), "baseUrl builds with native-transport");
}
}