1use std::cell::{Cell, RefCell};
11use std::collections::{BTreeMap, BTreeSet, HashMap, HashSet, VecDeque};
12
13use rusqlite::types::{ToSqlOutput, Value as SqlValue, ValueRef};
14use rusqlite::{Connection, OpenFlags, OptionalExtension};
15use serde::{Deserialize, Serialize};
16use serde_json::{Map, Value};
17use sha2::{Digest, Sha256};
18use ssp2::model::{Frame, MediaType, Message, MsgKind, Op, OpResult, PushStatus, SubStatus};
19use ssp2::primitives::RawJson;
20use ssp2::segment::{decode_rows_segment, Column, ColumnType, ColumnValue, Row, RowsSegment};
21use ssp2::{
22 decode_message, encode_message, encode_presence_publish, parse_control, ControlMessage,
23 PresenceKind,
24};
25
26use crate::api::{
27 ClientChangeBatch, ClientDiagnosticsHost, ClientDiagnosticsLease, ClientDiagnosticsReplica,
28 ClientDiagnosticsRequest, ClientDiagnosticsSchema, ClientDiagnosticsSnapshot,
29 ClientDiagnosticsStorage, ClientLimits, CommandEffects, CommitOperation,
30 CommitOperationOutcome, CommitOutcome, CommitOutcomeQuery, CommitOutcomeResolution,
31 CommitOutcomeStatus, ConflictRecord, CoverageSnapshot, DiagnosticLastChange,
32 DiagnosticLastRound, DiagnosticRoundCounters, DiagnosticSubscription, LeaseState,
33 LocalDataPurgeInput, LocalDataPurgeResult, LocalDataPurgeTarget, LocalDataRebootstrapInput,
34 LocalDataRebootstrapResult, Mutation, PresencePeer, QueryRow, QuerySnapshot, QueryValue,
35 RejectionDetails, RejectionRecord, ResolveCommitOutcomeInput, RowState, SchemaFloor,
36 SubscriptionStateView, SyncIntent, SyncOutcome, SyncReport, SyncStatusSnapshot, TableChange,
37 WindowBase, WindowChange, WindowCoverage, WindowState, WindowUnitRef,
38 CLIENT_DIAGNOSTICS_VERSION, MAX_DIAGNOSTIC_EXPECTED_SUBSCRIPTIONS,
39};
40use crate::schema::{parse_schema_json, ClientSchema, FtsIndexSchema, TableSchema};
41use crate::transport::{BlobDownload, BlobUploadGrant, SegmentRequest, Transport, TransportError};
42use crate::values::{
43 bytes_to_hex, canonical_scope_json, column_value_to_json, decode_row_bytes, encode_row_json,
44 json_to_column_value, json_to_scope_map, normalize_values_casing, render_row_id_json,
45 scope_map_to_json, sort_scope_map,
46};
47
48const DEFAULT_ACCEPT: u8 = 0b0111;
53const ACCEPT_INLINE_ROWS: u8 = 1 << 0;
54const ACCEPT_EXTERNAL_ROWS: u8 = 1 << 1;
55const ACCEPT_SQLITE: u8 = 1 << 2;
56const ACCEPT_SIGNED_URLS: u8 = 1 << 3;
57const MAX_DIAGNOSTIC_DOMAINS: usize = 256;
58
59const LOCAL_SCHEMA_VERSION_KEY: &str = "localSchemaVersion";
61const LOCAL_REVISION_KEY: &str = "localRevision";
62const CLIENT_ID_KEY: &str = "clientId";
63const LEASE_STATE_KEY: &str = "leaseState";
64const SCHEMA_FLOOR_KEY: &str = "schemaFloor";
65const SECURITY_PREFLIGHT_PENDING_KEY: &str = "securityPreflightPending";
70const LOCAL_REBOOTSTRAP_RECEIPT_VERSION: u8 = 2;
71const MAX_JS_SAFE_INTEGER: u64 = 9_007_199_254_740_991;
72const OUTBOX_INCOMPATIBLE_CODE: &str = "sync.outbox_incompatible";
75
76#[derive(Debug, Serialize, Deserialize)]
77#[serde(rename_all = "camelCase", deny_unknown_fields)]
78struct PersistedLocalDataRebootstrapReceipt {
79 version: u8,
80 retained_commits: u64,
81 reset_subscriptions: u64,
82}
83
84fn invalid_local_rebootstrap_receipt() -> String {
85 "sync.local_corrupt: persisted local rebootstrap receipt is invalid".to_owned()
86}
87
88fn encode_local_rebootstrap_receipt(
89 retained_commits: usize,
90 reset_subscriptions: usize,
91) -> Result<String, String> {
92 let retained_commits =
93 u64::try_from(retained_commits).map_err(|_| invalid_local_rebootstrap_receipt())?;
94 let reset_subscriptions =
95 u64::try_from(reset_subscriptions).map_err(|_| invalid_local_rebootstrap_receipt())?;
96 if retained_commits > MAX_JS_SAFE_INTEGER || reset_subscriptions > MAX_JS_SAFE_INTEGER {
97 return Err(invalid_local_rebootstrap_receipt());
98 }
99 serde_json::to_string(&PersistedLocalDataRebootstrapReceipt {
100 version: LOCAL_REBOOTSTRAP_RECEIPT_VERSION,
101 retained_commits,
102 reset_subscriptions,
103 })
104 .map_err(|_| invalid_local_rebootstrap_receipt())
105}
106
107fn decode_local_rebootstrap_receipt(value: &str) -> Result<(usize, usize), String> {
108 if value == "v1" {
110 return Ok((0, 0));
111 }
112 let receipt: PersistedLocalDataRebootstrapReceipt =
113 serde_json::from_str(value).map_err(|_| invalid_local_rebootstrap_receipt())?;
114 if receipt.version != LOCAL_REBOOTSTRAP_RECEIPT_VERSION
115 || receipt.retained_commits > MAX_JS_SAFE_INTEGER
116 || receipt.reset_subscriptions > MAX_JS_SAFE_INTEGER
117 {
118 return Err(invalid_local_rebootstrap_receipt());
119 }
120 Ok((
121 usize::try_from(receipt.retained_commits)
122 .map_err(|_| invalid_local_rebootstrap_receipt())?,
123 usize::try_from(receipt.reset_subscriptions)
124 .map_err(|_| invalid_local_rebootstrap_receipt())?,
125 ))
126}
127
128#[derive(Debug, Clone, Copy, PartialEq, Eq)]
129enum SubState {
130 Active,
131 Revoked,
132 Failed,
133}
134
135#[cfg(test)]
136mod observation_tests {
137 use super::*;
138 use crate::native_transport::HostTransport;
139 use serde_json::json;
140
141 fn client() -> SyncClient {
142 SyncClient::new(
143 "retry-test".to_owned(),
144 &json!({
145 "version": 1,
146 "tables": [{
147 "name": "tasks",
148 "primaryKey": "id",
149 "columns": [
150 { "name": "id", "type": "string", "nullable": false },
151 { "name": "project_id", "type": "string", "nullable": false }
152 ],
153 "scopes": [{ "pattern": "project:{project_id}" }]
154 }]
155 }),
156 ClientLimits::default(),
157 )
158 .expect("test client")
159 }
160
161 #[test]
162 fn background_retry_deadlines_back_off_and_reset() {
163 let mut client = client();
164 client.schedule_background_retry();
165 assert!(matches!(
166 client.drain_sync_intents().as_slice(),
167 [SyncIntent::Background { delay_ms: 250 }]
168 ));
169 client.schedule_background_retry();
170 assert!(matches!(
171 client.drain_sync_intents().as_slice(),
172 [SyncIntent::Background { delay_ms: 500 }]
173 ));
174 client.reset_background_retry();
175 client.schedule_background_retry();
176 assert!(matches!(
177 client.drain_sync_intents().as_slice(),
178 [SyncIntent::Background { delay_ms: 250 }]
179 ));
180 }
181
182 struct CountingRealtimeTransport {
183 connects: usize,
184 closes: usize,
185 }
186
187 impl Transport for CountingRealtimeTransport {
188 fn sync(&mut self, _request: &[u8]) -> Result<Vec<u8>, TransportError> {
189 Err(TransportError::new("sync.transport_failed", "offline"))
190 }
191
192 fn realtime_sync(&mut self, _request: &[u8]) -> Result<Vec<u8>, TransportError> {
193 Err(TransportError::new("sync.transport_failed", "offline"))
194 }
195
196 fn download_segment(
197 &mut self,
198 _request: &SegmentRequest,
199 ) -> Result<Vec<u8>, TransportError> {
200 Err(TransportError::new("sync.transport_failed", "offline"))
201 }
202
203 fn realtime_connect(&mut self) -> Result<(), TransportError> {
204 self.connects += 1;
205 Ok(())
206 }
207
208 fn realtime_send(&mut self, _text: &str) -> Result<(), TransportError> {
209 Ok(())
210 }
211
212 fn realtime_close(&mut self) -> Result<(), TransportError> {
213 self.closes += 1;
214 Ok(())
215 }
216 }
217
218 #[test]
219 fn realtime_connection_ownership_is_idempotent() {
220 let mut client = client();
221 let mut transport = CountingRealtimeTransport {
222 connects: 0,
223 closes: 0,
224 };
225 client
226 .connect_realtime(&mut transport)
227 .expect("first connect");
228 client
229 .connect_realtime(&mut transport)
230 .expect("idempotent connect");
231 assert_eq!(transport.connects, 1);
232 client.disconnect_realtime(&mut transport);
233 client.disconnect_realtime(&mut transport);
234 assert_eq!(transport.closes, 1);
235 client
236 .connect_realtime(&mut transport)
237 .expect("deliberate reconnect");
238 assert_eq!(transport.connects, 2);
239 }
240
241 #[test]
242 fn local_rebootstrap_is_atomic_idempotent_and_preserves_offline_work() {
243 let mut client = client();
244 client
245 .subscribe(
246 "repair-tasks".to_owned(),
247 "tasks".to_owned(),
248 vec![("project_id".to_owned(), vec!["p1".to_owned()])],
249 None,
250 )
251 .expect("subscribe");
252 {
253 let sub = client
254 .subs
255 .iter_mut()
256 .find(|sub| sub.id == "repair-tasks")
257 .expect("subscription");
258 sub.cursor = 42;
259 sub.synced_once = true;
260 let persisted = sub.clone();
261 client.persist_sub(&persisted);
262 }
263 for table in ["_syncular_base_tasks", "tasks"] {
264 client
265 .conn
266 .execute(
267 &format!(
268 "INSERT INTO {table}(id, project_id, _syncular_version) VALUES (?1, ?2, 1)"
269 ),
270 rusqlite::params!["server-row", "p1"],
271 )
272 .expect("seed server row");
273 }
274 let pending = client
275 .mutate(vec![Mutation::Upsert {
276 table: "tasks".to_owned(),
277 values: Map::from_iter([
278 ("id".to_owned(), Value::from("offline-row")),
279 ("project_id".to_owned(), Value::from("p1")),
280 ]),
281 base_version: None,
282 }])
283 .expect("queue offline work");
284 client.drain_change_batches();
285 client.drain_sync_intents();
286
287 assert_eq!(
288 client
289 .rebootstrap_local_data(&LocalDataRebootstrapInput {
290 rebootstrap_id: "support-case-001".to_owned(),
291 })
292 .expect("rebootstrap"),
293 LocalDataRebootstrapResult {
294 already_applied: false,
295 retained_commits: 1,
296 reset_subscriptions: 1,
297 }
298 );
299 let visible_ids = client
300 .conn
301 .prepare("SELECT id FROM tasks ORDER BY id")
302 .expect("prepare visible ids")
303 .query_map([], |row| row.get::<_, String>(0))
304 .expect("query visible ids")
305 .collect::<Result<Vec<_>, _>>()
306 .expect("collect visible ids");
307 assert_eq!(visible_ids, vec!["offline-row"]);
308 assert_eq!(client.pending_commit_ids(), vec![pending]);
309 assert_eq!(
310 client
311 .subscription_state("repair-tasks")
312 .expect("subscription")
313 .cursor,
314 -1
315 );
316 assert!(client.upgrading());
317 assert!(client.sync_needed());
318 assert!(matches!(
319 client.drain_sync_intents().as_slice(),
320 [SyncIntent::Interactive]
321 ));
322 assert_eq!(client.drain_change_batches().len(), 1);
323
324 assert_eq!(
325 client
326 .rebootstrap_local_data(&LocalDataRebootstrapInput {
327 rebootstrap_id: "support-case-001".to_owned(),
328 })
329 .expect("idempotent retry"),
330 LocalDataRebootstrapResult {
331 already_applied: true,
332 retained_commits: 1,
333 reset_subscriptions: 1,
334 }
335 );
336 assert!(client.drain_change_batches().is_empty());
337 }
338
339 #[test]
340 fn local_rebootstrap_receipt_codec_is_bounded_and_legacy_compatible() {
341 let encoded = encode_local_rebootstrap_receipt(3, 4).expect("encode receipt");
342 assert_eq!(
343 decode_local_rebootstrap_receipt(&encoded).expect("decode receipt"),
344 (3, 4)
345 );
346 assert_eq!(
347 decode_local_rebootstrap_receipt("v1").expect("legacy marker"),
348 (0, 0)
349 );
350 for malformed in [
351 "",
352 "{}",
353 r#"{"version":3,"retainedCommits":1,"resetSubscriptions":1}"#,
354 r#"{"version":2,"retainedCommits":1,"resetSubscriptions":1,"extra":true}"#,
355 r#"{"version":2,"retainedCommits":9007199254740992,"resetSubscriptions":1}"#,
356 ] {
357 assert_eq!(
358 decode_local_rebootstrap_receipt(malformed)
359 .expect_err("malformed receipt must fail"),
360 "sync.local_corrupt: persisted local rebootstrap receipt is invalid"
361 );
362 }
363 }
364
365 #[test]
366 fn local_rebootstrap_replays_the_original_receipt_after_reopen() {
367 let path = std::env::temp_dir().join(format!(
368 "syncular-rebootstrap-receipt-{}.db",
369 uuid::Uuid::new_v4()
370 ));
371 let schema = json!({
372 "version": 1,
373 "tables": [{
374 "name": "tasks",
375 "primaryKey": "id",
376 "columns": [
377 { "name": "id", "type": "string", "nullable": false },
378 { "name": "project_id", "type": "string", "nullable": false }
379 ],
380 "scopes": [{ "pattern": "project:{project_id}" }]
381 }]
382 });
383 let path_string = path.to_str().expect("UTF-8 temp path");
384
385 {
386 let mut first = SyncClient::open_path(
387 "repair-restart-client".to_owned(),
388 &schema,
389 ClientLimits::default(),
390 path_string,
391 )
392 .expect("first open");
393 first
394 .subscribe(
395 "repair-tasks".to_owned(),
396 "tasks".to_owned(),
397 vec![("project_id".to_owned(), vec!["p1".to_owned()])],
398 None,
399 )
400 .expect("subscribe");
401 first
402 .mutate(vec![Mutation::Upsert {
403 table: "tasks".to_owned(),
404 values: Map::from_iter([
405 ("id".to_owned(), Value::from("offline-row")),
406 ("project_id".to_owned(), Value::from("p1")),
407 ]),
408 base_version: None,
409 }])
410 .expect("queue offline work");
411 assert_eq!(
412 first
413 .rebootstrap_local_data(&LocalDataRebootstrapInput {
414 rebootstrap_id: "restart-receipt".to_owned(),
415 })
416 .expect("first rebootstrap"),
417 LocalDataRebootstrapResult {
418 already_applied: false,
419 retained_commits: 1,
420 reset_subscriptions: 1,
421 }
422 );
423 }
424
425 let mut reopened = SyncClient::open_path(
426 "repair-restart-client".to_owned(),
427 &schema,
428 ClientLimits::default(),
429 path_string,
430 )
431 .expect("reopen");
432 reopened.drain_change_batches();
433 reopened.drain_sync_intents();
434 assert_eq!(
435 reopened
436 .rebootstrap_local_data(&LocalDataRebootstrapInput {
437 rebootstrap_id: "restart-receipt".to_owned(),
438 })
439 .expect("receipt replay"),
440 LocalDataRebootstrapResult {
441 already_applied: true,
442 retained_commits: 1,
443 reset_subscriptions: 1,
444 }
445 );
446 assert!(reopened.drain_change_batches().is_empty());
447 assert!(reopened.drain_sync_intents().is_empty());
448 drop(reopened);
449 std::fs::remove_file(path).expect("remove temp database");
450 }
451
452 #[test]
453 fn local_rebootstrap_fails_closed_on_malformed_or_unreadable_receipts() {
454 let mut malformed = client();
455 malformed
456 .subscribe(
457 "repair-tasks".to_owned(),
458 "tasks".to_owned(),
459 vec![("project_id".to_owned(), vec!["p1".to_owned()])],
460 None,
461 )
462 .expect("subscribe");
463 malformed.set_meta("localRebootstrap:malformed", "{\"version\":2}");
464 malformed.drain_change_batches();
465 malformed.drain_sync_intents();
466 let malformed_error = malformed
467 .rebootstrap_local_data(&LocalDataRebootstrapInput {
468 rebootstrap_id: "malformed".to_owned(),
469 })
470 .expect_err("malformed receipt must fail closed");
471 assert_eq!(
472 malformed_error,
473 "sync.local_corrupt: persisted local rebootstrap receipt is invalid"
474 );
475 assert!(!malformed.upgrading());
476 assert_eq!(
477 malformed
478 .subscription_state("repair-tasks")
479 .expect("unchanged subscription")
480 .cursor,
481 -1
482 );
483 assert!(malformed.drain_change_batches().is_empty());
484 assert!(malformed.drain_sync_intents().is_empty());
485
486 let mut unreadable = client();
487 unreadable
488 .conn
489 .execute(
490 "INSERT INTO tasks(id, project_id, _syncular_version) VALUES (?1, ?2, 1)",
491 rusqlite::params!["server-row", "p1"],
492 )
493 .expect("seed visible row");
494 unreadable
495 .conn
496 .execute("DROP TABLE _syncular_meta", [])
497 .expect("break marker storage");
498 let unreadable_error = unreadable
499 .rebootstrap_local_data(&LocalDataRebootstrapInput {
500 rebootstrap_id: "unreadable".to_owned(),
501 })
502 .expect_err("unreadable marker storage must fail closed");
503 assert_eq!(
504 unreadable_error,
505 "sync.local_corrupt: persisted local rebootstrap receipt is unreadable"
506 );
507 let visible_rows = unreadable
508 .conn
509 .query_row("SELECT COUNT(*) FROM tasks", [], |row| row.get::<_, i64>(0))
510 .expect("visible projection remains");
511 assert_eq!(visible_rows, 1);
512 assert!(!unreadable.upgrading());
513 assert!(unreadable.drain_change_batches().is_empty());
514 assert!(unreadable.drain_sync_intents().is_empty());
515 }
516
517 #[test]
518 fn local_rebootstrap_cannot_bypass_schema_floor() {
519 let mut client = client();
520 client.set_schema_floor(Some(SchemaFloor {
521 required_schema_version: Some(2),
522 latest_schema_version: Some(2),
523 }));
524 let error = client
525 .rebootstrap_local_data(&LocalDataRebootstrapInput {
526 rebootstrap_id: "blocked-floor".to_owned(),
527 })
528 .expect_err("schema floor must block repair");
529 assert!(error.contains("cannot bypass an active schema-floor stop"));
530 }
531
532 #[test]
533 fn batched_push_acknowledgements_rebuild_overlay_once_per_response() {
534 let mut client = client();
535 const COMMIT_COUNT: usize = 32;
536
537 for index in 0..COMMIT_COUNT {
538 client
539 .mutate(vec![Mutation::Upsert {
540 table: "tasks".to_owned(),
541 values: Map::from_iter([
542 ("id".to_owned(), Value::from(format!("task-{index}"))),
543 ("project_id".to_owned(), Value::from("project-1")),
544 ]),
545 base_version: None,
546 }])
547 .expect("queue commit");
548 }
549
550 let (_, request_meta) = client.build_request(false);
551 assert_eq!(request_meta.pushed_ids.len(), COMMIT_COUNT);
552 let mut frames = vec![Frame::RespHeader {
553 required_schema_version: None,
554 latest_schema_version: None,
555 }];
556 frames.extend(request_meta.pushed_ids.iter().enumerate().map(
557 |(index, client_commit_id)| Frame::PushResult {
558 client_commit_id: client_commit_id.clone(),
559 status: PushStatus::Applied,
560 commit_seq: Some(index as i64 + 1),
561 results: vec![OpResult::Applied { op_index: 0 }],
562 },
563 ));
564 let response = Message {
565 msg_kind: MsgKind::Response,
566 frames,
567 };
568 let mut transport =
569 HostTransport::new_from_config(&json!({})).expect("no-network host transport");
570
571 client.overlay_rebuild_count.set(0);
572 client.outcome_prune_count.set(0);
573 let outcome = client.process_response(&mut transport, response, &request_meta);
574 assert!(matches!(outcome, SyncOutcome::Ok(_)));
575 assert!(client.pending_commit_ids().is_empty());
576 assert_eq!(
577 client.overlay_rebuild_count.get(),
578 1,
579 "one response must reconcile its acknowledged commits with one overlay rebuild"
580 );
581 assert_eq!(
582 client.outcome_prune_count.get(),
583 1,
584 "one response must enforce outcome retention once"
585 );
586 }
587
588 #[test]
589 fn mixed_push_results_reconcile_and_prune_once_per_response() {
590 let mut client = client();
591 for index in 0..4 {
592 client
593 .mutate(vec![Mutation::Upsert {
594 table: "tasks".to_owned(),
595 values: Map::from_iter([
596 ("id".to_owned(), Value::from(format!("task-{index}"))),
597 ("project_id".to_owned(), Value::from("project-1")),
598 ]),
599 base_version: None,
600 }])
601 .expect("queue commit");
602 }
603
604 let (_, request_meta) = client.build_request(false);
605 let ids = &request_meta.pushed_ids;
606 assert_eq!(ids.len(), 4);
607 let response = Message {
608 msg_kind: MsgKind::Response,
609 frames: vec![
610 Frame::RespHeader {
611 required_schema_version: None,
612 latest_schema_version: None,
613 },
614 Frame::PushResult {
615 client_commit_id: ids[0].clone(),
616 status: PushStatus::Applied,
617 commit_seq: Some(1),
618 results: vec![OpResult::Applied { op_index: 0 }],
619 },
620 Frame::PushResult {
621 client_commit_id: ids[1].clone(),
622 status: PushStatus::Cached,
623 commit_seq: Some(2),
624 results: vec![OpResult::Applied { op_index: 0 }],
625 },
626 Frame::PushResult {
627 client_commit_id: ids[2].clone(),
628 status: PushStatus::Rejected,
629 commit_seq: None,
630 results: vec![OpResult::Error {
631 op_index: 0,
632 code: "sync.validation_failed".to_owned(),
633 message: "rejected".to_owned(),
634 retryable: false,
635 }],
636 },
637 Frame::PushResult {
638 client_commit_id: ids[3].clone(),
639 status: PushStatus::Rejected,
640 commit_seq: None,
641 results: vec![OpResult::Error {
642 op_index: 0,
643 code: "sync.idempotency_cache_miss".to_owned(),
644 message: "retry".to_owned(),
645 retryable: true,
646 }],
647 },
648 ],
649 };
650 let mut transport =
651 HostTransport::new_from_config(&json!({})).expect("no-network host transport");
652
653 client.overlay_rebuild_count.set(0);
654 client.outcome_prune_count.set(0);
655 let outcome = client.process_response(&mut transport, response, &request_meta);
656 let SyncOutcome::Ok(report) = outcome else {
657 panic!("mixed push-result response failed");
658 };
659
660 assert_eq!(report.applied, ids[..2]);
661 assert_eq!(report.rejected, ids[2..3]);
662 assert_eq!(report.retryable, ids[3..4]);
663 assert_eq!(client.pending_commit_ids(), ids[3..4]);
664 assert_eq!(
665 client
666 .query("SELECT id FROM tasks ORDER BY id", &[])
667 .expect("query visible overlay"),
668 vec![Map::from_iter([("id".to_owned(), Value::from("task-3"))])]
669 );
670 assert_eq!(client.overlay_rebuild_count.get(), 1);
671 assert_eq!(client.outcome_prune_count.get(), 1);
672 }
673
674 #[test]
675 fn secondary_unique_collision_preserves_existing_synced_row() {
676 let client = SyncClient::new(
677 "unique-upsert-test".to_owned(),
678 &json!({
679 "version": 1,
680 "tables": [{
681 "name": "tasks",
682 "primaryKey": "id",
683 "columns": [
684 { "name": "id", "type": "string", "nullable": false },
685 { "name": "project_id", "type": "string", "nullable": false },
686 { "name": "title", "type": "string", "nullable": false }
687 ],
688 "scopes": [{ "pattern": "project:{project_id}" }],
689 "indexes": [{
690 "name": "tasks_by_project_title",
691 "columns": ["project_id", "title"],
692 "unique": true
693 }]
694 }]
695 }),
696 ClientLimits::default(),
697 )
698 .expect("test client");
699 let table = client.schema.table("tasks").expect("tasks table");
700 let sql = client.insert_row_sql(&base_table("tasks"), table);
701
702 client
703 .conn
704 .execute(&sql, rusqlite::params!["t1", "p1", "original", 1])
705 .expect("insert first row");
706 client
707 .conn
708 .execute(&sql, rusqlite::params!["t1", "p1", "updated", 2])
709 .expect("update same primary key");
710 client
711 .conn
712 .execute(&sql, rusqlite::params!["t2", "p1", "original", 1])
713 .expect("insert second row");
714 assert!(client
715 .conn
716 .execute(&sql, rusqlite::params!["t3", "p1", "original", 2])
717 .is_err());
718
719 let rows = client
720 .conn
721 .prepare("SELECT id, title FROM _syncular_base_tasks ORDER BY id")
722 .expect("prepare rows")
723 .query_map([], |row| {
724 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
725 })
726 .expect("query rows")
727 .collect::<Result<Vec<_>, _>>()
728 .expect("collect rows");
729 assert_eq!(
730 rows,
731 vec![
732 ("t1".to_owned(), "updated".to_owned()),
733 ("t2".to_owned(), "original".to_owned())
734 ]
735 );
736 }
737
738 #[test]
739 fn reopening_active_subscriptions_emits_a_catch_up_intent() {
740 let path = std::env::temp_dir().join(format!(
741 "syncular-startup-intent-{}.db",
742 uuid::Uuid::new_v4()
743 ));
744 let schema = json!({
745 "version": 1,
746 "tables": [{
747 "name": "tasks",
748 "primaryKey": "id",
749 "columns": [
750 { "name": "id", "type": "string", "nullable": false },
751 { "name": "project_id", "type": "string", "nullable": false }
752 ],
753 "scopes": [{ "pattern": "project:{project_id}" }]
754 }]
755 });
756
757 {
758 let mut first = SyncClient::open_path_with_identity(
759 None,
760 &schema,
761 ClientLimits::default(),
762 path.to_str().expect("UTF-8 temp path"),
763 )
764 .expect("first open");
765 first
766 .set_window(
767 &WindowBase {
768 table: "tasks".to_owned(),
769 variable: "project_id".to_owned(),
770 fixed_scopes: Vec::new(),
771 params: None,
772 },
773 &["persisted".to_owned()],
774 )
775 .expect("persist window");
776 }
777
778 let mut reopened = SyncClient::open_path_with_identity(
779 None,
780 &schema,
781 ClientLimits::default(),
782 path.to_str().expect("UTF-8 temp path"),
783 )
784 .expect("reopen");
785 assert!(reopened.sync_needed());
786 assert!(matches!(
787 reopened.drain_sync_intents().as_slice(),
788 [SyncIntent::Interactive]
789 ));
790 drop(reopened);
791 std::fs::remove_file(path).expect("remove temp database");
792 }
793
794 #[test]
795 fn reopening_preserves_immutable_subscription_identity_and_progress() {
796 let path = std::env::temp_dir().join(format!(
797 "syncular-subscription-identity-{}.db",
798 uuid::Uuid::new_v4()
799 ));
800 let schema = json!({
801 "version": 1,
802 "tables": [
803 {
804 "name": "tasks",
805 "primaryKey": "id",
806 "columns": [
807 { "name": "id", "type": "string", "nullable": false },
808 { "name": "project_id", "type": "string", "nullable": false }
809 ],
810 "scopes": [{ "pattern": "project:{project_id}" }]
811 },
812 {
813 "name": "docs",
814 "primaryKey": "id",
815 "columns": [
816 { "name": "id", "type": "string", "nullable": false },
817 { "name": "org_id", "type": "string", "nullable": false },
818 { "name": "project_id", "type": "string", "nullable": false }
819 ],
820 "scopes": [
821 { "pattern": "org:{org_id}" },
822 { "pattern": "project:{projectId}", "column": "project_id" }
823 ]
824 }
825 ]
826 });
827
828 {
829 let mut first = SyncClient::open_path_with_identity(
830 None,
831 &schema,
832 ClientLimits::default(),
833 path.to_str().expect("UTF-8 temp path"),
834 )
835 .expect("first open");
836 first
837 .subscribe(
838 "stable-subscription".to_owned(),
839 "tasks".to_owned(),
840 vec![(
841 "project_id".to_owned(),
842 vec!["p2".to_owned(), "p1".to_owned()],
843 )],
844 Some(r#"{"view":"v1"}"#.to_owned()),
845 )
846 .expect("persist subscription");
847 let persisted = {
848 let subscription = first
849 .subs
850 .iter_mut()
851 .find(|subscription| subscription.id == "stable-subscription")
852 .expect("subscription");
853 subscription.cursor = 41;
854 subscription.bootstrap_state = Some("resume-token".to_owned());
855 subscription.effective = Some(vec![(
856 "project_id".to_owned(),
857 vec!["p1".to_owned(), "p2".to_owned()],
858 )]);
859 subscription.synced_once = true;
860 subscription.clone()
861 };
862 first.persist_sub(&persisted);
863 }
864
865 let mut reopened = SyncClient::open_path_with_identity(
866 None,
867 &schema,
868 ClientLimits::default(),
869 path.to_str().expect("UTF-8 temp path"),
870 )
871 .expect("reopen");
872 let progress = reopened
873 .subscription_state("stable-subscription")
874 .expect("persisted state");
875 assert_eq!(progress.cursor, 41);
876 assert!(progress.has_resume_token);
877
878 reopened
879 .subscribe(
880 "stable-subscription".to_owned(),
881 "tasks".to_owned(),
882 vec![(
883 "project_id".to_owned(),
884 vec!["p1".to_owned(), "p2".to_owned(), "p1".to_owned()],
885 )],
886 Some(r#"{"view":"v1"}"#.to_owned()),
887 )
888 .expect("canonical intent is idempotent");
889 assert_eq!(
890 reopened
891 .subscription_state("stable-subscription")
892 .expect("unchanged state")
893 .cursor,
894 progress.cursor
895 );
896
897 for (table, scopes, params) in [
898 (
899 "tasks",
900 vec![("project_id".to_owned(), vec!["p1".to_owned()])],
901 Some(r#"{"view":"v1"}"#.to_owned()),
902 ),
903 (
904 "tasks",
905 vec![(
906 "project_id".to_owned(),
907 vec!["p2".to_owned(), "p1".to_owned()],
908 )],
909 Some(r#"{"view":"v2"}"#.to_owned()),
910 ),
911 (
912 "docs",
913 vec![
914 ("org_id".to_owned(), vec!["o1".to_owned()]),
915 ("projectId".to_owned(), vec!["p1".to_owned()]),
916 ],
917 Some(r#"{"view":"v1"}"#.to_owned()),
918 ),
919 ] {
920 let error = reopened
921 .subscribe(
922 "stable-subscription".to_owned(),
923 table.to_owned(),
924 scopes,
925 params,
926 )
927 .expect_err("identity rebind must fail");
928 assert!(error.starts_with("client.subscription_intent_mismatch:"));
929 assert_eq!(
930 reopened
931 .subscription_state("stable-subscription")
932 .expect("unchanged state")
933 .cursor,
934 progress.cursor
935 );
936 }
937
938 drop(reopened);
939 std::fs::remove_file(path).expect("remove temp database");
940 }
941
942 #[test]
943 fn reopening_clears_a_schema_floor_the_running_app_already_satisfies() {
944 let path = std::env::temp_dir().join(format!(
945 "syncular-satisfied-schema-floor-{}.db",
946 uuid::Uuid::new_v4()
947 ));
948 let schema = json!({
949 "version": 23,
950 "tables": [{
951 "name": "tasks",
952 "primaryKey": "id",
953 "columns": [
954 { "name": "id", "type": "string", "nullable": false },
955 { "name": "project_id", "type": "string", "nullable": false }
956 ],
957 "scopes": [{ "pattern": "project:{project_id}" }]
958 }]
959 });
960
961 {
962 let mut first = SyncClient::open_path_with_identity(
963 None,
964 &schema,
965 ClientLimits::default(),
966 path.to_str().expect("UTF-8 temp path"),
967 )
968 .expect("first open");
969 first
970 .subscribe(
971 "tasks".to_owned(),
972 "tasks".to_owned(),
973 vec![("project_id".to_owned(), vec!["p1".to_owned()])],
974 None,
975 )
976 .expect("persist subscription");
977 first.set_schema_floor(Some(SchemaFloor {
978 required_schema_version: Some(22),
979 latest_schema_version: Some(22),
980 }));
981 }
982
983 let mut reopened = SyncClient::open_path_with_identity(
984 None,
985 &schema,
986 ClientLimits::default(),
987 path.to_str().expect("UTF-8 temp path"),
988 )
989 .expect("reopen");
990 assert!(reopened.schema_floor().is_none());
991 assert!(reopened.get_meta(SCHEMA_FLOOR_KEY).is_none());
992 assert!(reopened.sync_needed());
993 assert!(matches!(
994 reopened.drain_sync_intents().as_slice(),
995 [SyncIntent::Interactive]
996 ));
997 drop(reopened);
998 std::fs::remove_file(path).expect("remove temp database");
999 }
1000
1001 #[test]
1002 fn reopening_keeps_an_unsatisfied_schema_floor_stopped() {
1003 let path = std::env::temp_dir().join(format!(
1004 "syncular-unsatisfied-schema-floor-{}.db",
1005 uuid::Uuid::new_v4()
1006 ));
1007 let schema = json!({
1008 "version": 1,
1009 "tables": [{
1010 "name": "tasks",
1011 "primaryKey": "id",
1012 "columns": [
1013 { "name": "id", "type": "string", "nullable": false },
1014 { "name": "project_id", "type": "string", "nullable": false }
1015 ],
1016 "scopes": [{ "pattern": "project:{project_id}" }]
1017 }]
1018 });
1019
1020 {
1021 let mut first = SyncClient::open_path_with_identity(
1022 None,
1023 &schema,
1024 ClientLimits::default(),
1025 path.to_str().expect("UTF-8 temp path"),
1026 )
1027 .expect("first open");
1028 first
1029 .subscribe(
1030 "tasks".to_owned(),
1031 "tasks".to_owned(),
1032 vec![("project_id".to_owned(), vec!["p1".to_owned()])],
1033 None,
1034 )
1035 .expect("persist subscription");
1036 first.set_schema_floor(Some(SchemaFloor {
1037 required_schema_version: Some(2),
1038 latest_schema_version: Some(2),
1039 }));
1040 }
1041
1042 let reopened = SyncClient::open_path_with_identity(
1043 None,
1044 &schema,
1045 ClientLimits::default(),
1046 path.to_str().expect("UTF-8 temp path"),
1047 )
1048 .expect("reopen");
1049 assert_eq!(
1050 reopened.schema_floor(),
1051 Some(&SchemaFloor {
1052 required_schema_version: Some(2),
1053 latest_schema_version: Some(2),
1054 })
1055 );
1056 assert!(!reopened.sync_needed());
1057 drop(reopened);
1058 std::fs::remove_file(path).expect("remove temp database");
1059 }
1060
1061 #[test]
1062 fn schema_bump_precedes_index_ddl_and_prunes_removed_subscriptions() {
1063 let path = std::env::temp_dir().join(format!(
1064 "syncular-indexed-column-bump-{}.db",
1065 uuid::Uuid::new_v4()
1066 ));
1067 let old_schema = json!({
1068 "version": 1,
1069 "tables": [
1070 {
1071 "name": "tasks",
1072 "primaryKey": "id",
1073 "columns": [
1074 { "name": "id", "type": "string", "nullable": false },
1075 { "name": "project_id", "type": "string", "nullable": false }
1076 ],
1077 "scopes": [{ "pattern": "project:{project_id}" }]
1078 },
1079 {
1080 "name": "legacy",
1081 "primaryKey": "id",
1082 "columns": [
1083 { "name": "id", "type": "string", "nullable": false },
1084 { "name": "project_id", "type": "string", "nullable": false }
1085 ],
1086 "scopes": [{ "pattern": "project:{project_id}" }]
1087 }
1088 ]
1089 });
1090 {
1091 let mut first = SyncClient::open_path_with_identity(
1092 None,
1093 &old_schema,
1094 ClientLimits::default(),
1095 path.to_str().expect("UTF-8 temp path"),
1096 )
1097 .expect("open old schema");
1098 first
1099 .subscribe(
1100 "legacy-sub".to_owned(),
1101 "legacy".to_owned(),
1102 vec![("project_id".to_owned(), vec!["p1".to_owned()])],
1103 None,
1104 )
1105 .expect("persist legacy subscription");
1106 }
1107
1108 let new_schema = json!({
1109 "version": 2,
1110 "tables": [{
1111 "name": "tasks",
1112 "primaryKey": "id",
1113 "columns": [
1114 { "name": "id", "type": "string", "nullable": false },
1115 { "name": "project_id", "type": "string", "nullable": false },
1116 { "name": "facility_membership_id", "type": "string", "nullable": true }
1117 ],
1118 "scopes": [{ "pattern": "project:{project_id}" }],
1119 "indexes": [{
1120 "name": "tasks_by_membership",
1121 "columns": ["project_id", "facility_membership_id"],
1122 "unique": false
1123 }]
1124 }]
1125 });
1126 let upgraded = SyncClient::open_path_with_identity(
1127 None,
1128 &new_schema,
1129 ClientLimits::default(),
1130 path.to_str().expect("UTF-8 temp path"),
1131 )
1132 .expect("open upgraded schema");
1133 let columns = upgraded
1134 .conn
1135 .prepare("PRAGMA table_info(tasks)")
1136 .expect("prepare columns")
1137 .query_map([], |row| row.get::<_, String>(1))
1138 .expect("query columns")
1139 .collect::<Result<Vec<_>, _>>()
1140 .expect("collect columns");
1141 assert!(columns.contains(&"facility_membership_id".to_owned()));
1142 let index_count: i64 = upgraded
1143 .conn
1144 .query_row(
1145 "SELECT COUNT(*) FROM sqlite_master WHERE type = 'index' AND name = 'tasks_by_membership'",
1146 [],
1147 |row| row.get(0),
1148 )
1149 .expect("query index");
1150 assert_eq!(index_count, 1);
1151 assert!(upgraded.subscription_state("legacy-sub").is_none());
1152 drop(upgraded);
1153 std::fs::remove_file(path).expect("remove temp database");
1154 }
1155
1156 #[test]
1157 fn migrates_pre_envelope_outcome_journal_additively() {
1158 let path = std::env::temp_dir().join(format!(
1159 "syncular-outcome-migration-{}.db",
1160 uuid::Uuid::new_v4()
1161 ));
1162 let conn = Connection::open(&path).expect("open legacy database");
1163 conn.execute_batch(
1164 "CREATE TABLE _syncular_commit_outcomes (
1165 seq INTEGER PRIMARY KEY AUTOINCREMENT,
1166 client_commit_id TEXT NOT NULL UNIQUE,
1167 status TEXT NOT NULL,
1168 recorded_at_ms INTEGER NOT NULL,
1169 results_json TEXT NOT NULL,
1170 resolution TEXT NOT NULL DEFAULT 'active',
1171 resolved_at_ms INTEGER,
1172 replacement_client_commit_id TEXT);",
1173 )
1174 .expect("create legacy outcome journal");
1175 drop(conn);
1176 let schema = json!({
1177 "version": 1,
1178 "tables": [{
1179 "name": "tasks",
1180 "primaryKey": "id",
1181 "columns": [
1182 { "name": "id", "type": "string", "nullable": false },
1183 { "name": "project_id", "type": "string", "nullable": false }
1184 ],
1185 "scopes": [{ "pattern": "project:{project_id}" }]
1186 }]
1187 });
1188 let client = SyncClient::open_path(
1189 "migration-native".to_owned(),
1190 &schema,
1191 ClientLimits::default(),
1192 path.to_str().expect("UTF-8 temp path"),
1193 )
1194 .expect("migrate database");
1195 let has_operations = client
1196 .conn
1197 .prepare("PRAGMA table_info(_syncular_commit_outcomes)")
1198 .expect("prepare table info")
1199 .query_map([], |row| row.get::<_, String>(1))
1200 .expect("query table info")
1201 .filter_map(Result::ok)
1202 .any(|column| column == "operations_json");
1203 assert!(has_operations);
1204 drop(client);
1205 std::fs::remove_file(path).expect("remove temp database");
1206 }
1207
1208 #[test]
1209 fn durable_conflict_outcome_and_resolution_survive_reopen() {
1210 let path = std::env::temp_dir().join(format!(
1211 "syncular-durable-outcome-{}.db",
1212 uuid::Uuid::new_v4()
1213 ));
1214 let schema = json!({
1215 "version": 1,
1216 "tables": [{
1217 "name": "tasks",
1218 "primaryKey": "id",
1219 "columns": [
1220 { "name": "id", "type": "string", "nullable": false },
1221 { "name": "project_id", "type": "string", "nullable": false }
1222 ],
1223 "scopes": [{ "pattern": "project:{project_id}" }]
1224 }]
1225 });
1226 let conflict = ConflictRecord {
1227 client_commit_id: "losing-commit".to_owned(),
1228 op_index: 0,
1229 table: "tasks".to_owned(),
1230 row_id: "t1".to_owned(),
1231 code: "sync.version_conflict".to_owned(),
1232 message: "stale base version".to_owned(),
1233 server_version: 2,
1234 server_row: Map::from_iter([("id".to_owned(), json!("t1"))]),
1235 operation: Some(CommitOperation {
1236 table: "tasks".to_owned(),
1237 row_id: "t1".to_owned(),
1238 op: "upsert".to_owned(),
1239 base_version: Some(1),
1240 values: None,
1241 changed_fields: None,
1242 }),
1243 };
1244 let failed_operations = vec![
1245 OutboxOp {
1246 upsert: true,
1247 table: "tasks".to_owned(),
1248 row_id: "t1".to_owned(),
1249 base_version: Some(1),
1250 values: None,
1251 changed_fields: None,
1252 },
1253 OutboxOp {
1254 upsert: true,
1255 table: "tasks".to_owned(),
1256 row_id: "status-event-1".to_owned(),
1257 base_version: Some(0),
1258 values: None,
1259 changed_fields: None,
1260 },
1261 ];
1262
1263 {
1264 let mut first = SyncClient::open_path(
1265 "durable-native".to_owned(),
1266 &schema,
1267 ClientLimits::default(),
1268 path.to_str().expect("UTF-8 temp path"),
1269 )
1270 .expect("first open");
1271 first
1272 .begin_observation("test_outcome")
1273 .expect("begin outcome");
1274 first
1275 .persist_commit_outcome(
1276 "losing-commit",
1277 CommitOutcomeStatus::Conflict,
1278 &[CommitOperationOutcome::Conflict {
1279 conflict: conflict.clone(),
1280 }],
1281 Some(&failed_operations),
1282 )
1283 .expect("persist outcome");
1284 first.conflicts.push(conflict);
1285 first
1286 .finish_observation(
1287 "test_outcome",
1288 ChangeAccumulator {
1289 conflicts: true,
1290 outcomes: true,
1291 ..ChangeAccumulator::default()
1292 },
1293 )
1294 .expect("commit outcome");
1295 }
1296
1297 {
1298 let mut reopened = SyncClient::open_path(
1299 "durable-native".to_owned(),
1300 &schema,
1301 ClientLimits::default(),
1302 path.to_str().expect("UTF-8 temp path"),
1303 )
1304 .expect("reopen");
1305 assert_eq!(reopened.conflicts().len(), 1);
1306 let outcome = reopened
1307 .commit_outcome("losing-commit")
1308 .expect("read outcome")
1309 .expect("outcome");
1310 assert_eq!(outcome.status, CommitOutcomeStatus::Conflict);
1311 let operations = outcome.operations.expect("aggregate envelope");
1312 assert_eq!(operations.len(), 2);
1313 assert_eq!(operations[1].row_id, "status-event-1");
1314 let resolved = reopened
1315 .resolve_commit_outcome(ResolveCommitOutcomeInput {
1316 client_commit_id: "losing-commit".to_owned(),
1317 resolution: CommitOutcomeResolution::ResolvedKeepServer,
1318 replacement_client_commit_id: None,
1319 })
1320 .expect("resolve");
1321 assert_eq!(
1322 resolved.resolution,
1323 CommitOutcomeResolution::ResolvedKeepServer
1324 );
1325 assert!(reopened.conflicts().is_empty());
1326 }
1327
1328 let reopened = SyncClient::open_path(
1329 "durable-native".to_owned(),
1330 &schema,
1331 ClientLimits::default(),
1332 path.to_str().expect("UTF-8 temp path"),
1333 )
1334 .expect("second reopen");
1335 assert!(reopened.conflicts().is_empty());
1336 assert_eq!(
1337 reopened
1338 .commit_outcome("losing-commit")
1339 .expect("read outcome")
1340 .expect("outcome")
1341 .resolution,
1342 CommitOutcomeResolution::ResolvedKeepServer
1343 );
1344 drop(reopened);
1345 std::fs::remove_file(path).expect("remove temp database");
1346 }
1347
1348 #[test]
1349 fn file_snapshot_reader_matches_owner_rows_revision_and_coverage() {
1350 let path =
1351 std::env::temp_dir().join(format!("syncular-read-sidecar-{}.db", uuid::Uuid::new_v4()));
1352 let schema = json!({
1353 "version": 1,
1354 "tables": [{
1355 "name": "tasks",
1356 "primaryKey": "id",
1357 "columns": [
1358 { "name": "id", "type": "string", "nullable": false },
1359 { "name": "project_id", "type": "string", "nullable": false }
1360 ],
1361 "scopes": [{ "pattern": "project:{project_id}" }]
1362 }]
1363 });
1364 let mut client = SyncClient::open_path_with_identity(
1365 Some("sidecar-client".to_owned()),
1366 &schema,
1367 ClientLimits::default(),
1368 path.to_str().expect("UTF-8 temp path"),
1369 )
1370 .expect("open owner");
1371 let base = WindowBase {
1372 table: "tasks".to_owned(),
1373 variable: "project_id".to_owned(),
1374 fixed_scopes: Vec::new(),
1375 params: None,
1376 };
1377 client
1378 .set_window(&base, &["one".to_owned()])
1379 .expect("set window");
1380 client
1381 .mutate(vec![Mutation::Upsert {
1382 table: "tasks".to_owned(),
1383 values: Map::from_iter([
1384 ("id".to_owned(), Value::from("t1")),
1385 ("project_id".to_owned(), Value::from("one")),
1386 ]),
1387 base_version: None,
1388 }])
1389 .expect("local mutate");
1390
1391 let coverage = [WindowCoverage {
1392 base,
1393 units: vec!["one".to_owned(), "missing".to_owned()],
1394 }];
1395 let owner = client
1396 .query_snapshot(
1397 "SELECT id, project_id, _sync_version AS server_version FROM tasks ORDER BY id",
1398 &[],
1399 &coverage,
1400 )
1401 .expect("owner snapshot");
1402 let mut reader = FileQuerySnapshotReader::new(path.to_string_lossy());
1403 let sidecar = reader
1404 .query_snapshot(
1405 "SELECT id, project_id, _sync_version AS server_version FROM tasks ORDER BY id",
1406 &[],
1407 &coverage,
1408 )
1409 .expect("sidecar snapshot");
1410
1411 assert_eq!(sidecar.revision, owner.revision);
1412 assert_eq!(sidecar.rows, owner.rows);
1413 assert_eq!(
1414 serde_json::to_value(&sidecar.coverage).expect("serialize sidecar coverage"),
1415 serde_json::to_value(&owner.coverage).expect("serialize owner coverage")
1416 );
1417 assert_eq!(sidecar.revision, "2");
1418 assert_eq!(sidecar.rows[0]["id"], "t1");
1419 assert_eq!(sidecar.rows[0]["server_version"], -1);
1420 assert!(!sidecar.coverage.complete);
1421 assert_eq!(sidecar.coverage.pending.len(), 1);
1422 assert_eq!(sidecar.coverage.missing.len(), 1);
1423
1424 drop(reader);
1425 drop(client);
1426 std::fs::remove_file(path).expect("remove temp database");
1427 }
1428
1429 #[test]
1430 fn local_fts_projection_tracks_optimistic_overlay_rebuilds() {
1431 let schema = json!({
1432 "version": 1,
1433 "tables": [{
1434 "name": "catalogue_codes",
1435 "primaryKey": "id",
1436 "columns": [
1437 { "name": "id", "type": "string", "nullable": false },
1438 { "name": "release_id", "type": "string", "nullable": false },
1439 { "name": "code", "type": "string", "nullable": false },
1440 { "name": "title", "type": "string", "nullable": false }
1441 ],
1442 "scopes": [{ "pattern": "release:{release_id}" }],
1443 "ftsIndexes": [{
1444 "name": "catalogue_codes_fts",
1445 "columns": ["code", "title"],
1446 "tokenize": "unicode61 remove_diacritics 2"
1447 }]
1448 }]
1449 });
1450 let mut client = SyncClient::new("fts-test".to_owned(), &schema, ClientLimits::default())
1451 .expect("FTS5 client");
1452 let insert_trigger: String = client
1453 .conn
1454 .query_row(
1455 "SELECT sql FROM sqlite_master WHERE type='trigger' AND name='catalogue_codes_fts_ai'",
1456 [],
1457 |row| row.get(0),
1458 )
1459 .expect("insert trigger");
1460 let replace_guard: String = client
1461 .conn
1462 .query_row(
1463 "SELECT sql FROM sqlite_master WHERE type='trigger' AND name='catalogue_codes_fts_bi'",
1464 [],
1465 |row| row.get(0),
1466 )
1467 .expect("replace guard");
1468 assert!(!insert_trigger.contains("DELETE FROM"));
1469 assert!(replace_guard.contains("BEFORE INSERT"));
1470 assert!(replace_guard.contains("WHEN EXISTS"));
1471 let search = |client: &SyncClient, query: &str| {
1472 client
1473 .query(
1474 "SELECT c.id FROM catalogue_codes_fts f JOIN catalogue_codes c ON CAST(c.id AS TEXT) = f._syncular_source_id WHERE catalogue_codes_fts MATCH ?1 ORDER BY c.id",
1475 &[Value::from(query)],
1476 )
1477 .expect("FTS query")
1478 };
1479
1480 client
1481 .mutate(vec![Mutation::Upsert {
1482 table: "catalogue_codes".to_owned(),
1483 values: Map::from_iter([
1484 ("id".to_owned(), Value::from("c1")),
1485 ("release_id".to_owned(), Value::from("r1")),
1486 ("code".to_owned(), Value::from("A01")),
1487 ("title".to_owned(), Value::from("Cholera")),
1488 ]),
1489 base_version: None,
1490 }])
1491 .expect("insert code");
1492 assert_eq!(search(&client, "cholera").len(), 1);
1493
1494 client
1495 .mutate(vec![Mutation::Upsert {
1496 table: "catalogue_codes".to_owned(),
1497 values: Map::from_iter([
1498 ("id".to_owned(), Value::from("c1")),
1499 ("release_id".to_owned(), Value::from("r1")),
1500 ("code".to_owned(), Value::from("A01")),
1501 ("title".to_owned(), Value::from("Enteric infection")),
1502 ]),
1503 base_version: None,
1504 }])
1505 .expect("update code");
1506 assert!(search(&client, "cholera").is_empty());
1507 assert_eq!(search(&client, "enteric").len(), 1);
1508
1509 client
1510 .mutate(vec![Mutation::Delete {
1511 table: "catalogue_codes".to_owned(),
1512 row_id: "c1".to_owned(),
1513 base_version: None,
1514 }])
1515 .expect("delete code");
1516 assert!(search(&client, "enteric").is_empty());
1517 }
1518
1519 #[test]
1520 fn application_authorized_local_purge_is_exact_atomic_and_idempotent() {
1521 let schema = json!({
1522 "version": 1,
1523 "tables": [{
1524 "name": "patient_notes",
1525 "primaryKey": "id",
1526 "columns": [
1527 { "name": "id", "type": "string", "nullable": false },
1528 { "name": "practice_id", "type": "string", "nullable": false },
1529 { "name": "encryption_key_id", "type": "string", "nullable": false },
1530 { "name": "title", "type": "string", "nullable": false }
1531 ],
1532 "scopes": [{ "pattern": "practice:{practice_id}" }],
1533 "ftsIndexes": [{
1534 "name": "patient_notes_fts",
1535 "columns": ["title"],
1536 "tokenize": "unicode61 remove_diacritics 2"
1537 }]
1538 }]
1539 });
1540 let mut client = SyncClient::new(
1541 "local-purge-test".to_owned(),
1542 &schema,
1543 ClientLimits::default(),
1544 )
1545 .expect("local purge client");
1546 client
1547 .conn
1548 .execute(
1549 "INSERT INTO _syncular_base_patient_notes(id, practice_id, encryption_key_id, title, _syncular_version) VALUES
1550 ('target', 'practice-1', 'key-revoked', 'Target original', 1),
1551 ('unrelated', 'practice-1', 'key-held', 'Unrelated original', 1)",
1552 [],
1553 )
1554 .expect("seed base rows");
1555 client.overlay_dirty.set(true);
1556 client.rebuild_overlay_if_dirty();
1557
1558 let note = |id: &str, key_id: &str, title: &str| {
1559 Map::from_iter([
1560 ("id".to_owned(), Value::from(id)),
1561 ("practice_id".to_owned(), Value::from("practice-1")),
1562 ("encryption_key_id".to_owned(), Value::from(key_id)),
1563 ("title".to_owned(), Value::from(title)),
1564 ])
1565 };
1566 let doomed = client
1567 .mutate(vec![
1568 Mutation::Upsert {
1569 table: "patient_notes".to_owned(),
1570 values: note("target", "key-revoked", "Target changed"),
1571 base_version: None,
1572 },
1573 Mutation::Upsert {
1574 table: "patient_notes".to_owned(),
1575 values: note("unrelated", "key-held", "Unrelated changed"),
1576 base_version: None,
1577 },
1578 ])
1579 .expect("doomed commit");
1580 let kept = client
1581 .mutate(vec![Mutation::Upsert {
1582 table: "patient_notes".to_owned(),
1583 values: note("kept", "key-held", "Kept optimistic"),
1584 base_version: None,
1585 }])
1586 .expect("kept commit");
1587 client.drain_change_batches();
1588
1589 let input = LocalDataPurgeInput {
1590 purge_id: "purge-001".to_owned(),
1591 targets: vec![LocalDataPurgeTarget {
1592 table: "patient_notes".to_owned(),
1593 selectors: BTreeMap::from([(
1594 "encryption_key_id".to_owned(),
1595 vec!["key-revoked".to_owned()],
1596 )]),
1597 }],
1598 };
1599 assert_eq!(
1600 client.purge_local_data(&input).expect("apply purge"),
1601 LocalDataPurgeResult {
1602 already_applied: false,
1603 purged_rows: 1,
1604 dropped_commits: 1,
1605 }
1606 );
1607 let rows = client
1608 .query("SELECT id, title FROM patient_notes ORDER BY id", &[])
1609 .expect("visible rows");
1610 assert_eq!(rows.len(), 2);
1611 assert_eq!(rows[0]["id"], "kept");
1612 assert_eq!(rows[1]["id"], "unrelated");
1613 assert_eq!(rows[1]["title"], "Unrelated original");
1614 let fts = client
1615 .query(
1616 "SELECT n.id FROM patient_notes_fts f JOIN patient_notes n ON CAST(n.id AS TEXT) = f._syncular_source_id WHERE patient_notes_fts MATCH 'target'",
1617 &[],
1618 )
1619 .expect("fts query");
1620 assert!(fts.is_empty());
1621 assert_eq!(client.outbox.len(), 1);
1622 assert_eq!(client.outbox[0].client_commit_id, kept);
1623 let outcome = client
1624 .commit_outcome(&doomed)
1625 .expect("read doomed outcome")
1626 .expect("doomed outcome");
1627 assert_eq!(outcome.status, CommitOutcomeStatus::Rejected);
1628 match &outcome.results[0] {
1629 CommitOperationOutcome::Error { rejection } => {
1630 assert_eq!(rejection.code, "client.local_data_purged");
1631 assert_eq!(rejection.client_commit_id, doomed);
1632 }
1633 other => panic!("expected local purge rejection, got {other:?}"),
1634 }
1635 assert_eq!(client.drain_change_batches().len(), 1);
1636 assert_eq!(
1637 client.purge_local_data(&input).expect("retry purge"),
1638 LocalDataPurgeResult {
1639 already_applied: true,
1640 purged_rows: 0,
1641 dropped_commits: 0,
1642 }
1643 );
1644 let conflicting = LocalDataPurgeInput {
1645 purge_id: input.purge_id.clone(),
1646 targets: vec![LocalDataPurgeTarget {
1647 table: "patient_notes".to_owned(),
1648 selectors: BTreeMap::from([(
1649 "encryption_key_id".to_owned(),
1650 vec!["key-held".to_owned()],
1651 )]),
1652 }],
1653 };
1654 assert!(client
1655 .purge_local_data(&conflicting)
1656 .expect_err("id collision must fail")
1657 .contains("already used with a different plan"));
1658 }
1659}
1660
1661impl SubState {
1662 fn name(self) -> &'static str {
1663 match self {
1664 SubState::Active => "active",
1665 SubState::Revoked => "revoked",
1666 SubState::Failed => "failed",
1667 }
1668 }
1669
1670 fn parse(value: &str) -> Self {
1671 match value {
1672 "revoked" => Self::Revoked,
1673 "failed" => Self::Failed,
1674 _ => Self::Active,
1675 }
1676 }
1677}
1678
1679#[derive(Debug, Clone)]
1680struct Subscription {
1681 id: String,
1682 table: String,
1683 requested: Vec<(String, Vec<String>)>,
1684 params: Option<String>,
1685 cursor: i64,
1686 bootstrap_state: Option<String>,
1688 state: SubState,
1689 reason_code: Option<String>,
1690 effective: Option<Vec<(String, Vec<String>)>>,
1693 synced_once: bool,
1694}
1695
1696#[derive(Debug, Clone)]
1697struct OutboxOp {
1698 upsert: bool,
1699 table: String,
1700 row_id: String,
1701 base_version: Option<i64>,
1702 values: Option<Map<String, Value>>,
1705 changed_fields: Option<Vec<String>>,
1707}
1708
1709impl From<&OutboxOp> for CommitOperation {
1710 fn from(operation: &OutboxOp) -> Self {
1711 Self {
1712 table: operation.table.clone(),
1713 row_id: operation.row_id.clone(),
1714 op: if operation.upsert { "upsert" } else { "delete" }.to_owned(),
1715 base_version: operation.base_version,
1716 values: operation.values.clone(),
1717 changed_fields: operation.changed_fields.clone(),
1718 }
1719 }
1720}
1721
1722#[derive(Debug, Clone)]
1723struct OutboxCommit {
1724 client_commit_id: String,
1725 ops: Vec<OutboxOp>,
1726}
1727
1728#[derive(Debug, Clone)]
1729struct CompiledLocalDataPurgeTarget {
1730 table: String,
1731 selectors: Vec<(String, Vec<String>)>,
1732}
1733
1734struct StoredCommitOutcomeRow {
1735 sequence: i64,
1736 client_commit_id: String,
1737 status: String,
1738 recorded_at_ms: i64,
1739 results_json: String,
1740 operations_json: Option<String>,
1741 resolution: String,
1742 resolved_at_ms: Option<i64>,
1743 replacement_client_commit_id: Option<String>,
1744}
1745
1746enum SectionError {
1749 FailClosed,
1750 Abort(String, String),
1751}
1752
1753struct RequestMeta {
1754 pushed_ids: Vec<String>,
1755 fresh: Vec<(String, bool)>,
1758 accept: u8,
1759 deferred_commits: usize,
1763}
1764
1765#[derive(Default)]
1766struct ChangeAccumulator {
1767 tables: BTreeMap<String, Option<BTreeSet<String>>>,
1769 windows: BTreeMap<(String, String), BTreeSet<String>>,
1770 status: bool,
1771 conflicts: bool,
1772 rejections: bool,
1773 outcomes: bool,
1774}
1775
1776impl ChangeAccumulator {
1777 fn table(&mut self, table: &str) {
1778 self.tables.insert(table.to_owned(), None);
1779 }
1780
1781 fn scope(&mut self, table: &str, key: String) {
1782 match self.tables.get_mut(table) {
1783 Some(None) => {}
1784 Some(Some(keys)) => {
1785 keys.insert(key);
1786 }
1787 None => {
1788 self.tables
1789 .insert(table.to_owned(), Some(BTreeSet::from([key])));
1790 }
1791 }
1792 }
1793
1794 fn window(&mut self, base_key: &str, table: &str, unit: &str) {
1795 self.windows
1796 .entry((base_key.to_owned(), table.to_owned()))
1797 .or_default()
1798 .insert(unit.to_owned());
1799 }
1800
1801 fn touched(&self) -> bool {
1802 !self.tables.is_empty()
1803 || !self.windows.is_empty()
1804 || self.status
1805 || self.conflicts
1806 || self.rejections
1807 || self.outcomes
1808 }
1809}
1810
1811const PUSH_OPS_PER_REQUEST: usize = 500;
1817const MAX_LOCAL_PURGE_TARGETS: usize = 64;
1818const MAX_LOCAL_PURGE_SELECTORS: usize = 8;
1819const MAX_LOCAL_PURGE_VALUES: usize = 128;
1820const MAX_LOCAL_PURGE_VALUE_LENGTH: usize = 256;
1821pub const SECURITY_PREFLIGHT_REQUIRED_CODE: &str = "client.security_preflight_required";
1822
1823pub struct SyncClient {
1824 conn: Connection,
1825 schema: ClientSchema,
1826 client_id: String,
1827 limits: ClientLimits,
1828 subs: Vec<Subscription>,
1829 outbox: Vec<OutboxCommit>,
1830 conflicts: Vec<ConflictRecord>,
1831 rejections: Vec<RejectionRecord>,
1832 schema_floor: Option<SchemaFloor>,
1833 lease_state: Option<LeaseState>,
1835 stopped: bool,
1837 upgrading: bool,
1839 sync_needed: bool,
1841 realtime_connected: bool,
1842 presence: HashMap<String, HashMap<String, PresencePeer>>,
1844 now_ms: Option<i64>,
1847 encryption: crate::values::EncryptionConfig,
1851 security_preflight: bool,
1854 insert_sql: RefCell<HashMap<String, String>>,
1860 overlay_dirty: Cell<bool>,
1864 #[cfg(test)]
1867 overlay_rebuild_count: Cell<usize>,
1868 #[cfg(test)]
1871 outcome_prune_count: Cell<usize>,
1872 change_queue: VecDeque<ClientChangeBatch>,
1874 sync_intent_queue: VecDeque<SyncIntent>,
1875 retry_delay_ms: u64,
1877 last_round: Option<DiagnosticLastRound>,
1878 last_change: Option<DiagnosticLastChange>,
1879}
1880
1881fn quote_ident(name: &str) -> String {
1882 format!("\"{}\"", name.replace('"', "\"\""))
1883}
1884
1885fn is_local_operation_code_like(value: &str) -> bool {
1886 let bytes = value.as_bytes();
1887 bytes.first().is_some_and(u8::is_ascii_alphanumeric)
1888 && bytes
1889 .iter()
1890 .all(|byte| byte.is_ascii_alphanumeric() || matches!(*byte, b'.' | b'_' | b':' | b'-'))
1891}
1892
1893fn base_table(name: &str) -> String {
1894 quote_ident(&format!("_syncular_base_{name}"))
1895}
1896
1897type PendingEvict = (String, String, Vec<(String, Vec<String>)>);
1902
1903fn window_base_key(base: &WindowBase) -> String {
1904 format!(
1905 "{}\0{}\0{}",
1906 base.table,
1907 base.variable,
1908 canonical_scope_json(&base.fixed_scopes)
1909 )
1910}
1911
1912fn unit_scopes(base: &WindowBase, unit: &str) -> Vec<(String, Vec<String>)> {
1914 let mut scopes = base.fixed_scopes.clone();
1915 scopes.retain(|(k, _)| k != &base.variable);
1916 scopes.push((base.variable.clone(), vec![unit.to_owned()]));
1917 scopes
1918}
1919
1920fn derive_sub_id(base: &WindowBase, unit: &str) -> String {
1924 let canonical = canonical_scope_json(&unit_scopes(base, unit));
1925 let digest = Sha256::digest(canonical.as_bytes());
1926 let hex = bytes_to_hex(&digest);
1927 format!("w:{}:{}", base.table, &hex[..16])
1928}
1929
1930fn visible_table(name: &str) -> String {
1931 quote_ident(name)
1932}
1933
1934const FTS_SOURCE_ID_COLUMN: &str = "_syncular_source_id";
1935
1936fn is_synced_table_name(name: &str) -> bool {
1939 if name.starts_with("sqlite_") {
1940 return false;
1941 }
1942 if name.starts_with("_syncular_base_") {
1943 return true; }
1945 !name.starts_with("_syncular_")
1947}
1948
1949fn blob_id_for(bytes: &[u8]) -> String {
1951 let digest = Sha256::digest(bytes);
1952 format!("sha256:{}", bytes_to_hex(&digest))
1953}
1954
1955enum RowParam<'a> {
1959 Cell(&'a Option<ColumnValue>),
1960 Version(i64),
1961}
1962
1963impl rusqlite::ToSql for RowParam<'_> {
1964 fn to_sql(&self) -> rusqlite::Result<ToSqlOutput<'_>> {
1965 Ok(match self {
1966 RowParam::Version(v) => ToSqlOutput::Owned(SqlValue::Integer(*v)),
1967 RowParam::Cell(cell) => match cell {
1968 None => ToSqlOutput::Owned(SqlValue::Null),
1969 Some(ColumnValue::String(s)) => ToSqlOutput::Borrowed(ValueRef::Text(s.as_bytes())),
1970 Some(ColumnValue::Integer(i)) => ToSqlOutput::Owned(SqlValue::Integer(*i)),
1971 Some(ColumnValue::Float(f)) => ToSqlOutput::Owned(SqlValue::Real(*f)),
1972 Some(ColumnValue::Boolean(b)) => {
1973 ToSqlOutput::Owned(SqlValue::Integer(i64::from(*b)))
1974 }
1975 Some(ColumnValue::Json(raw)) | Some(ColumnValue::BlobRef(raw)) => {
1976 ToSqlOutput::Borrowed(ValueRef::Text(raw.0.as_bytes()))
1977 }
1978 Some(ColumnValue::Bytes(b)) | Some(ColumnValue::Crdt(b)) => {
1980 ToSqlOutput::Borrowed(ValueRef::Blob(b))
1981 }
1982 },
1983 })
1984 }
1985}
1986
1987fn image_cell_param<'a>(column: &Column, value: ValueRef<'a>) -> Result<ToSqlOutput<'a>, String> {
1995 use ssp2::segment::ColumnType;
1996 let mismatch = || {
1997 Err(format!(
1998 "image column {:?} holds a value of the wrong type",
1999 column.name
2000 ))
2001 };
2002 match value {
2003 ValueRef::Null => {
2004 if !column.nullable {
2005 return Err(format!(
2006 "image column {:?} is NULL but not nullable",
2007 column.name
2008 ));
2009 }
2010 Ok(ToSqlOutput::Owned(SqlValue::Null))
2011 }
2012 ValueRef::Integer(i) => match column.ty {
2013 ColumnType::Integer => Ok(ToSqlOutput::Borrowed(value)),
2014 ColumnType::Boolean => Ok(ToSqlOutput::Owned(SqlValue::Integer(i64::from(i != 0)))),
2015 ColumnType::Float => Ok(ToSqlOutput::Owned(SqlValue::Real(i as f64))),
2016 _ => mismatch(),
2017 },
2018 ValueRef::Real(_) => match column.ty {
2019 ColumnType::Float => Ok(ToSqlOutput::Borrowed(value)),
2020 _ => mismatch(),
2021 },
2022 ValueRef::Text(t) => {
2023 std::str::from_utf8(t)
2024 .map_err(|_| format!("image column {:?} is not UTF-8", column.name))?;
2025 match column.ty {
2026 ColumnType::String | ColumnType::Json | ColumnType::BlobRef => {
2027 Ok(ToSqlOutput::Borrowed(value))
2028 }
2029 _ => mismatch(),
2030 }
2031 }
2032 ValueRef::Blob(_) => match column.ty {
2033 ColumnType::Bytes | ColumnType::Crdt => Ok(ToSqlOutput::Borrowed(value)),
2035 _ => mismatch(),
2036 },
2037 }
2038}
2039
2040fn sql_ref_to_json(column: &Column, value: rusqlite::types::ValueRef<'_>) -> Value {
2041 use rusqlite::types::ValueRef;
2042 match value {
2043 ValueRef::Null => Value::Null,
2044 ValueRef::Integer(i) => match column.ty {
2045 ssp2::segment::ColumnType::Boolean => Value::Bool(i != 0),
2046 ssp2::segment::ColumnType::Float => {
2047 serde_json::Number::from_f64(i as f64).map_or(Value::Null, Value::Number)
2048 }
2049 _ => Value::from(i),
2050 },
2051 ValueRef::Real(f) => serde_json::Number::from_f64(f).map_or(Value::Null, Value::Number),
2052 ValueRef::Text(t) => Value::from(String::from_utf8_lossy(t).into_owned()),
2053 ValueRef::Blob(b) => {
2054 let mut map = Map::new();
2055 map.insert("$bytes".to_owned(), Value::from(bytes_to_hex(b)));
2056 Value::Object(map)
2057 }
2058 }
2059}
2060
2061fn json_param_to_sql(value: &Value) -> Result<SqlValue, String> {
2065 Ok(match value {
2066 Value::Null => SqlValue::Null,
2067 Value::Bool(b) => SqlValue::Integer(i64::from(*b)),
2068 Value::Number(n) => {
2069 if let Some(i) = n.as_i64() {
2070 SqlValue::Integer(i)
2071 } else if let Some(f) = n.as_f64() {
2072 SqlValue::Real(f)
2073 } else {
2074 return Err(format!("query param number {n} is out of range"));
2075 }
2076 }
2077 Value::String(s) => SqlValue::Text(s.clone()),
2078 Value::Object(_) => {
2079 if let Some(hex) = value.get("$bytes").and_then(Value::as_str) {
2080 SqlValue::Blob(crate::values::hex_to_bytes(hex)?)
2081 } else if let Some(decimal) = value.get("$bigint").and_then(Value::as_str) {
2082 SqlValue::Integer(decimal.parse::<i64>().map_err(|_| {
2083 format!("query bigint param {decimal:?} is outside SQLite's i64 range")
2084 })?)
2085 } else {
2086 return Err(
2087 "query object param must be a {$bytes: hex} or {$bigint: decimal} value"
2088 .to_owned(),
2089 );
2090 }
2091 }
2092 Value::Array(_) => return Err("query array params are not supported".to_owned()),
2093 })
2094}
2095
2096fn sql_ref_to_json_dynamic(value: rusqlite::types::ValueRef<'_>) -> Value {
2101 use rusqlite::types::ValueRef;
2102 match value {
2103 ValueRef::Null => Value::Null,
2104 ValueRef::Integer(i) => {
2105 const MAX_SAFE_INTEGER: i64 = 9_007_199_254_740_991;
2109 if (-MAX_SAFE_INTEGER..=MAX_SAFE_INTEGER).contains(&i) {
2110 Value::from(i)
2111 } else {
2112 let mut map = Map::new();
2113 map.insert("$bigint".to_owned(), Value::from(i.to_string()));
2114 Value::Object(map)
2115 }
2116 }
2117 ValueRef::Real(f) => serde_json::Number::from_f64(f).map_or(Value::Null, Value::Number),
2118 ValueRef::Text(t) => Value::from(String::from_utf8_lossy(t).into_owned()),
2119 ValueRef::Blob(b) => {
2120 let mut map = Map::new();
2121 map.insert("$bytes".to_owned(), Value::from(bytes_to_hex(b)));
2122 Value::Object(map)
2123 }
2124 }
2125}
2126
2127fn query_connection(
2128 conn: &Connection,
2129 sql: &str,
2130 params: &[QueryValue],
2131) -> Result<Vec<QueryRow>, String> {
2132 crate::query_guard::assert_read_only_query(sql)?;
2133 let lowered_sql = crate::query_guard::lower_public_query_sql(sql);
2134 let bound: Vec<SqlValue> = params
2135 .iter()
2136 .map(json_param_to_sql)
2137 .collect::<Result<_, _>>()?;
2138 let mut stmt = conn.prepare(&lowered_sql).map_err(|e| e.to_string())?;
2139 let column_names: Vec<String> = stmt.column_names().into_iter().map(str::to_owned).collect();
2140 let bound_refs: Vec<&dyn rusqlite::ToSql> =
2141 bound.iter().map(|v| v as &dyn rusqlite::ToSql).collect();
2142 let mut sql_rows = stmt
2143 .query(bound_refs.as_slice())
2144 .map_err(|e| e.to_string())?;
2145 let mut out = Vec::new();
2146 while let Some(row) = sql_rows.next().map_err(|e| e.to_string())? {
2147 let mut record = Map::new();
2148 for (i, name) in column_names.iter().enumerate() {
2149 let value = row.get_ref(i).map_err(|e| e.to_string())?;
2150 record.insert(name.clone(), sql_ref_to_json_dynamic(value));
2151 }
2152 out.push(record);
2153 }
2154 Ok(out)
2155}
2156
2157fn persisted_window_state(conn: &Connection, base: &WindowBase) -> Result<WindowState, String> {
2158 let mut stmt = conn
2159 .prepare(
2160 "SELECT windows.unit, subscriptions.state_json
2161 FROM _syncular_windows AS windows
2162 JOIN _syncular_subscriptions AS subscriptions
2163 ON subscriptions.id = windows.sub_id
2164 WHERE windows.base = ?1
2165 ORDER BY windows.unit ASC",
2166 )
2167 .map_err(|error| error.to_string())?;
2168 let rows = stmt
2169 .query_map(rusqlite::params![window_base_key(base)], |row| {
2170 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
2171 })
2172 .map_err(|error| error.to_string())?;
2173 let mut units = Vec::new();
2174 let mut pending = Vec::new();
2175 for row in rows {
2176 let (unit, raw) = row.map_err(|error| error.to_string())?;
2177 let state: Value = serde_json::from_str(&raw)
2178 .map_err(|error| format!("invalid persisted window subscription: {error}"))?;
2179 let is_pending = state.get("status").and_then(Value::as_str) != Some("active")
2180 || state.get("cursor").and_then(Value::as_i64).unwrap_or(-1) < 0
2181 || state
2182 .get("bootstrapState")
2183 .is_some_and(|value| !value.is_null());
2184 if is_pending {
2185 pending.push(unit.clone());
2186 }
2187 units.push(unit);
2188 }
2189 Ok(WindowState { units, pending })
2190}
2191
2192fn snapshot_connection(
2193 conn: &Connection,
2194 sql: &str,
2195 params: &[Value],
2196 coverage: &[WindowCoverage],
2197) -> Result<QuerySnapshot, String> {
2198 conn.execute_batch("SAVEPOINT syncular_snapshot_read")
2199 .map_err(|error| error.to_string())?;
2200 let result = (|| {
2201 let revision = conn
2202 .query_row(
2203 "SELECT value FROM _syncular_meta WHERE key = ?1",
2204 rusqlite::params![LOCAL_REVISION_KEY],
2205 |row| row.get::<_, String>(0),
2206 )
2207 .ok()
2208 .and_then(|value| value.parse::<u64>().ok())
2209 .unwrap_or(0);
2210 let rows = query_connection(conn, sql, params)?;
2211 let mut pending = Vec::new();
2212 let mut missing = Vec::new();
2213 for requested in coverage {
2214 let base_key = window_base_key(&requested.base);
2215 let state = persisted_window_state(conn, &requested.base)?;
2216 for unit in BTreeSet::from_iter(requested.units.iter().cloned()) {
2217 let reference = WindowUnitRef {
2218 base_key: base_key.clone(),
2219 unit: unit.clone(),
2220 };
2221 if !state.units.iter().any(|held| held == &unit) {
2222 missing.push(reference);
2223 } else if state.pending.iter().any(|held| held == &unit) {
2224 pending.push(reference);
2225 }
2226 }
2227 }
2228 Ok(QuerySnapshot {
2229 revision: revision.to_string(),
2230 rows,
2231 coverage: CoverageSnapshot {
2232 complete: pending.is_empty() && missing.is_empty(),
2233 pending,
2234 missing,
2235 },
2236 })
2237 })();
2238 match result {
2239 Ok(snapshot) => {
2240 conn.execute_batch("RELEASE syncular_snapshot_read")
2241 .map_err(|error| error.to_string())?;
2242 Ok(snapshot)
2243 }
2244 Err(error) => {
2245 let _ = conn.execute_batch(
2246 "ROLLBACK TO syncular_snapshot_read; RELEASE syncular_snapshot_read",
2247 );
2248 Err(error)
2249 }
2250 }
2251}
2252
2253pub struct FileQuerySnapshotReader {
2258 path: String,
2259 conn: Option<Connection>,
2260}
2261
2262impl FileQuerySnapshotReader {
2263 #[must_use]
2264 pub fn new(path: impl Into<String>) -> Self {
2265 Self {
2266 path: path.into(),
2267 conn: None,
2268 }
2269 }
2270
2271 fn connection(&mut self) -> Result<&Connection, String> {
2272 if self.conn.is_none() {
2273 let conn = Connection::open_with_flags(
2274 &self.path,
2275 OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_NO_MUTEX,
2276 )
2277 .map_err(|error| format!("open read sidecar {:?}: {error}", self.path))?;
2278 conn.busy_timeout(std::time::Duration::from_millis(250))
2279 .map_err(|error| error.to_string())?;
2280 self.conn = Some(conn);
2281 }
2282 self.conn
2283 .as_ref()
2284 .ok_or_else(|| "read sidecar connection missing".to_owned())
2285 }
2286
2287 pub fn query_snapshot(
2288 &mut self,
2289 sql: &str,
2290 params: &[Value],
2291 coverage: &[WindowCoverage],
2292 ) -> Result<QuerySnapshot, String> {
2293 snapshot_connection(self.connection()?, sql, params, coverage)
2294 }
2295}
2296
2297impl SyncClient {
2298 pub fn new_with_identity(
2299 client_id: Option<String>,
2300 schema_json: &Value,
2301 limits: ClientLimits,
2302 ) -> Result<Self, String> {
2303 let resolved = client_id.unwrap_or_else(|| uuid::Uuid::new_v4().to_string());
2304 let conn = Connection::open_in_memory().map_err(|e| e.to_string())?;
2305 Self::with_connection(resolved, schema_json, limits, conn)
2306 }
2307
2308 pub fn new(
2309 client_id: String,
2310 schema_json: &Value,
2311 limits: ClientLimits,
2312 ) -> Result<Self, String> {
2313 let conn = Connection::open_in_memory().map_err(|e| e.to_string())?;
2314 Self::with_connection(client_id, schema_json, limits, conn)
2315 }
2316
2317 pub fn open_path(
2323 client_id: String,
2324 schema_json: &Value,
2325 limits: ClientLimits,
2326 path: &str,
2327 ) -> Result<Self, String> {
2328 let conn = Connection::open(path).map_err(|e| format!("open db {path:?}: {e}"))?;
2329 Self::with_connection(client_id, schema_json, limits, conn)
2330 }
2331
2332 pub fn open_path_with_identity(
2333 client_id: Option<String>,
2334 schema_json: &Value,
2335 limits: ClientLimits,
2336 path: &str,
2337 ) -> Result<Self, String> {
2338 let conn = Connection::open(path).map_err(|e| format!("open db {path:?}: {e}"))?;
2339 conn.busy_timeout(std::time::Duration::from_millis(250))
2345 .map_err(|error| format!("configure db {path:?} busy timeout: {error}"))?;
2346 conn.pragma_update(None, "journal_mode", "WAL")
2347 .map_err(|error| format!("configure db {path:?} WAL mode: {error}"))?;
2348 let persisted = conn
2349 .query_row(
2350 "SELECT value FROM _syncular_meta WHERE key = 'clientId'",
2351 [],
2352 |row| row.get::<_, String>(0),
2353 )
2354 .ok();
2355 let resolved = persisted
2356 .clone()
2357 .or(client_id.clone())
2358 .unwrap_or_else(|| uuid::Uuid::new_v4().to_string());
2359 if let (Some(existing), Some(requested)) = (persisted, client_id) {
2360 if existing != requested {
2361 return Err(format!(
2362 "client.identity_mismatch: this database belongs to {existing:?}; refusing to rebind it to {requested:?}"
2363 ));
2364 }
2365 }
2366 Self::with_connection(resolved, schema_json, limits, conn)
2367 }
2368
2369 #[must_use]
2370 pub fn client_id(&self) -> &str {
2371 &self.client_id
2372 }
2373
2374 pub fn with_connection(
2381 client_id: String,
2382 schema_json: &Value,
2383 limits: ClientLimits,
2384 conn: Connection,
2385 ) -> Result<Self, String> {
2386 if limits.outcome_retention_max_entries == Some(0) {
2387 return Err(
2388 "sync.invalid_request: outcomeRetentionMaxEntries must be positive".to_owned(),
2389 );
2390 }
2391 let schema = parse_schema_json(schema_json)?;
2392 let mut client = SyncClient {
2393 conn,
2394 schema,
2395 client_id,
2396 limits,
2397 subs: Vec::new(),
2398 outbox: Vec::new(),
2399 conflicts: Vec::new(),
2400 rejections: Vec::new(),
2401 schema_floor: None,
2402 lease_state: None,
2403 stopped: false,
2404 upgrading: false,
2405 sync_needed: false,
2406 realtime_connected: false,
2407 presence: HashMap::new(),
2408 now_ms: None,
2409 encryption: crate::values::EncryptionConfig::default(),
2410 security_preflight: false,
2411 insert_sql: RefCell::new(HashMap::new()),
2412 overlay_dirty: Cell::new(false),
2413 #[cfg(test)]
2414 overlay_rebuild_count: Cell::new(0),
2415 #[cfg(test)]
2416 outcome_prune_count: Cell::new(0),
2417 change_queue: VecDeque::new(),
2418 sync_intent_queue: VecDeque::new(),
2419 retry_delay_ms: 250,
2420 last_round: None,
2421 last_change: None,
2422 };
2423 client
2427 .conn
2428 .set_prepared_statement_cache_capacity(64.max(client.schema.tables.len() * 4));
2429 client.create_bookkeeping_tables()?;
2433 match client.get_meta(CLIENT_ID_KEY) {
2434 Some(existing) if existing != client.client_id => {
2435 return Err(format!(
2436 "client.identity_mismatch: this database belongs to {existing:?}; refusing to rebind it to {:?}",
2437 client.client_id
2438 ));
2439 }
2440 None => client.set_meta(CLIENT_ID_KEY, &client.client_id),
2441 _ => {}
2442 }
2443 client.restore_persisted_state()?;
2444 let marker = client
2445 .get_meta(LOCAL_SCHEMA_VERSION_KEY)
2446 .and_then(|value| value.parse::<i32>().ok());
2447 match marker {
2448 None => {
2449 client.create_synced_tables()?;
2450 client.set_meta(LOCAL_SCHEMA_VERSION_KEY, &client.schema.version.to_string());
2451 }
2452 Some(version) if version == client.schema.version => {
2453 client.create_synced_tables()?;
2454 }
2455 Some(_) => client.run_schema_reset()?,
2456 }
2457 client.clear_satisfied_persisted_schema_floor();
2458 client.prune_unknown_subscriptions()?;
2459 if marker == Some(client.schema.version) && !client.outbox.is_empty() {
2460 client.overlay_dirty.set(true);
2463 client.rebuild_overlay();
2464 }
2465 client.enqueue_startup_sync_if_needed();
2471 Ok(client)
2472 }
2473
2474 pub fn set_now_ms(&mut self, now_ms: i64) {
2477 self.now_ms = Some(now_ms);
2478 }
2479
2480 pub fn set_encryption(&mut self, encryption: crate::values::EncryptionConfig) {
2484 self.encryption = encryption;
2485 }
2486
2487 #[must_use]
2488 pub fn security_lifecycle(&self) -> &'static str {
2489 if self.security_preflight {
2490 "preflight"
2491 } else {
2492 "active"
2493 }
2494 }
2495
2496 #[must_use]
2497 pub fn security_preflight(&self) -> bool {
2498 self.security_preflight
2499 }
2500
2501 pub fn begin_security_preflight(&mut self) {
2508 self.seal_security_on_teardown();
2509 self.set_meta(SECURITY_PREFLIGHT_PENDING_KEY, "1");
2512 }
2513
2514 pub fn seal_security_on_teardown(&mut self) {
2524 self.security_preflight = true;
2525 self.encryption = crate::values::EncryptionConfig::default();
2526 self.sync_intent_queue.clear();
2527 }
2528
2529 pub fn activate_security(
2531 &mut self,
2532 encryption: crate::values::EncryptionConfig,
2533 ) -> Result<(), String> {
2534 if !self.security_preflight {
2535 return Err(
2536 "sync.invalid_request: activateSecurity requires security preflight".to_owned(),
2537 );
2538 }
2539 self.encryption = encryption;
2540 self.security_preflight = false;
2541 self.delete_meta(SECURITY_PREFLIGHT_PENDING_KEY);
2542 self.enqueue_startup_sync_if_needed();
2543 Ok(())
2544 }
2545
2546 fn clock_now_ms(&self) -> i64 {
2547 self.now_ms.unwrap_or_else(|| {
2548 std::time::SystemTime::now()
2549 .duration_since(std::time::UNIX_EPOCH)
2550 .map(|d| d.as_millis() as i64)
2551 .unwrap_or(0)
2552 })
2553 }
2554
2555 fn create_bookkeeping_tables(&self) -> Result<(), String> {
2556 self.conn
2558 .execute_batch(
2559 "CREATE TABLE IF NOT EXISTS _syncular_outbox (
2560 seq INTEGER PRIMARY KEY AUTOINCREMENT,
2561 commit_id TEXT NOT NULL UNIQUE, ops_json TEXT NOT NULL);
2562 CREATE TABLE IF NOT EXISTS _syncular_commit_outcomes (
2563 seq INTEGER PRIMARY KEY AUTOINCREMENT,
2564 client_commit_id TEXT NOT NULL UNIQUE,
2565 status TEXT NOT NULL CHECK(status IN ('applied', 'cached', 'conflict', 'rejected')),
2566 recorded_at_ms INTEGER NOT NULL,
2567 results_json TEXT NOT NULL,
2568 operations_json TEXT,
2569 resolution TEXT NOT NULL DEFAULT 'active'
2570 CHECK(resolution IN ('active', 'resolved_keep_server', 'superseded', 'dismissed')),
2571 resolved_at_ms INTEGER,
2572 replacement_client_commit_id TEXT);
2573 CREATE INDEX IF NOT EXISTS _syncular_commit_outcomes_resolution_seq
2574 ON _syncular_commit_outcomes(resolution, seq);
2575 CREATE TABLE IF NOT EXISTS _syncular_subscriptions (
2576 id TEXT PRIMARY KEY, tbl TEXT NOT NULL, state_json TEXT NOT NULL);
2577 CREATE TABLE IF NOT EXISTS _syncular_meta (
2578 key TEXT PRIMARY KEY, value TEXT NOT NULL);
2579 CREATE TABLE IF NOT EXISTS _syncular_windows (
2580 base TEXT NOT NULL, unit TEXT NOT NULL, sub_id TEXT NOT NULL,
2581 PRIMARY KEY (base, unit));
2582 CREATE TABLE IF NOT EXISTS _syncular_window_pending_evict (
2583 sub_id TEXT PRIMARY KEY, tbl TEXT NOT NULL,
2584 effective_scopes TEXT NOT NULL);",
2585 )
2586 .map_err(|e| e.to_string())?;
2587 let _ = self
2590 .conn
2591 .execute_batch("ALTER TABLE _syncular_commit_outcomes ADD COLUMN operations_json TEXT");
2592 if self.get_meta(LOCAL_REVISION_KEY).is_none() {
2593 self.set_meta(LOCAL_REVISION_KEY, "0");
2594 }
2595 if self.schema_has_blobs() {
2600 self.conn
2601 .execute_batch(
2602 "CREATE TABLE IF NOT EXISTS _syncular_blobs (blob_id TEXT PRIMARY KEY,
2603 bytes BLOB NOT NULL, byte_length INTEGER NOT NULL,
2604 media_type TEXT, refcount INTEGER NOT NULL DEFAULT 0,
2605 created_at_ms INTEGER NOT NULL,
2606 last_used_ms INTEGER NOT NULL DEFAULT 0);
2607 CREATE TABLE IF NOT EXISTS _syncular_blob_uploads (blob_id TEXT PRIMARY KEY,
2608 media_type TEXT, created_at_ms INTEGER NOT NULL);",
2609 )
2610 .map_err(|e| e.to_string())?;
2611 let _ = self.conn.execute_batch(
2615 "ALTER TABLE _syncular_blobs ADD COLUMN last_used_ms INTEGER NOT NULL DEFAULT 0",
2616 );
2617 }
2618 Ok(())
2619 }
2620
2621 fn schema_has_blobs(&self) -> bool {
2623 self.schema
2624 .tables
2625 .iter()
2626 .any(|t| t.columns.iter().any(|c| c.ty == ColumnType::BlobRef))
2627 }
2628
2629 fn get_meta(&self, key: &str) -> Option<String> {
2632 self.conn
2633 .query_row(
2634 "SELECT value FROM _syncular_meta WHERE key = ?1",
2635 rusqlite::params![key],
2636 |row| row.get::<_, String>(0),
2637 )
2638 .ok()
2639 }
2640
2641 fn get_meta_strict(&self, key: &str) -> Result<Option<String>, String> {
2642 self.conn
2643 .query_row(
2644 "SELECT value FROM _syncular_meta WHERE key = ?1",
2645 rusqlite::params![key],
2646 |row| row.get::<_, String>(0),
2647 )
2648 .optional()
2649 .map_err(|_| {
2650 "sync.local_corrupt: persisted local rebootstrap receipt is unreadable".to_owned()
2651 })
2652 }
2653
2654 fn set_meta(&self, key: &str, value: &str) {
2655 let _ = self.conn.execute(
2656 "INSERT OR REPLACE INTO _syncular_meta (key, value) VALUES (?1, ?2)",
2657 rusqlite::params![key, value],
2658 );
2659 }
2660
2661 fn delete_meta(&self, key: &str) {
2662 let _ = self.conn.execute(
2663 "DELETE FROM _syncular_meta WHERE key = ?1",
2664 rusqlite::params![key],
2665 );
2666 }
2667
2668 fn restore_persisted_state(&mut self) -> Result<(), String> {
2669 self.subs = {
2670 let mut stmt = self
2671 .conn
2672 .prepare("SELECT id, tbl, state_json FROM _syncular_subscriptions ORDER BY id ASC")
2673 .map_err(|error| error.to_string())?;
2674 let rows = stmt
2675 .query_map([], |row| {
2676 Ok((
2677 row.get::<_, String>(0)?,
2678 row.get::<_, String>(1)?,
2679 row.get::<_, String>(2)?,
2680 ))
2681 })
2682 .map_err(|error| error.to_string())?;
2683 let mut subscriptions = Vec::new();
2684 for row in rows {
2685 let (id, table, raw) = row.map_err(|error| error.to_string())?;
2686 let state: Value = serde_json::from_str(&raw)
2687 .map_err(|error| format!("invalid persisted subscription {id:?}: {error}"))?;
2688 let requested = json_to_scope_map(
2689 state.get("requested").unwrap_or(&Value::Object(Map::new())),
2690 )?;
2691 let effective = state
2692 .get("effectiveScopes")
2693 .filter(|value| !value.is_null())
2694 .map(json_to_scope_map)
2695 .transpose()?;
2696 subscriptions.push(Subscription {
2697 id,
2698 table,
2699 requested,
2700 params: state
2701 .get("params")
2702 .and_then(Value::as_str)
2703 .map(str::to_owned),
2704 cursor: state.get("cursor").and_then(Value::as_i64).unwrap_or(-1),
2705 bootstrap_state: state
2706 .get("bootstrapState")
2707 .and_then(Value::as_str)
2708 .map(str::to_owned),
2709 state: SubState::parse(
2710 state
2711 .get("status")
2712 .and_then(Value::as_str)
2713 .unwrap_or("active"),
2714 ),
2715 reason_code: state
2716 .get("reasonCode")
2717 .and_then(Value::as_str)
2718 .map(str::to_owned),
2719 effective,
2720 synced_once: state
2721 .get("syncedOnce")
2722 .and_then(Value::as_bool)
2723 .unwrap_or(false),
2724 });
2725 }
2726 subscriptions
2727 };
2728
2729 self.outbox = {
2730 let mut stmt = self
2731 .conn
2732 .prepare("SELECT commit_id, ops_json FROM _syncular_outbox ORDER BY seq ASC")
2733 .map_err(|error| error.to_string())?;
2734 let rows = stmt
2735 .query_map([], |row| {
2736 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
2737 })
2738 .map_err(|error| error.to_string())?;
2739 let mut commits = Vec::new();
2740 for row in rows {
2741 let (client_commit_id, raw) = row.map_err(|error| error.to_string())?;
2742 let entries: Vec<Value> = serde_json::from_str(&raw).map_err(|error| {
2743 format!("invalid persisted outbox {client_commit_id:?}: {error}")
2744 })?;
2745 let mut ops = Vec::with_capacity(entries.len());
2746 for entry in entries {
2747 let op = entry.get("op").and_then(Value::as_str).unwrap_or("delete");
2748 ops.push(OutboxOp {
2749 upsert: op == "upsert",
2750 table: entry
2751 .get("table")
2752 .and_then(Value::as_str)
2753 .ok_or_else(|| "persisted outbox operation missing table".to_owned())?
2754 .to_owned(),
2755 row_id: entry
2756 .get("rowId")
2757 .and_then(Value::as_str)
2758 .ok_or_else(|| "persisted outbox operation missing rowId".to_owned())?
2759 .to_owned(),
2760 base_version: entry.get("baseVersion").and_then(Value::as_i64),
2761 values: entry.get("values").and_then(Value::as_object).cloned(),
2762 changed_fields: entry.get("changedFields").and_then(Value::as_array).map(
2763 |values| {
2764 values
2765 .iter()
2766 .filter_map(Value::as_str)
2767 .map(str::to_owned)
2768 .collect()
2769 },
2770 ),
2771 });
2772 }
2773 commits.push(OutboxCommit {
2774 client_commit_id,
2775 ops,
2776 });
2777 }
2778 commits
2779 };
2780
2781 self.prune_commit_outcomes()?;
2782 let active = self.commit_outcomes(CommitOutcomeQuery {
2783 active_only: true,
2784 ..CommitOutcomeQuery::default()
2785 })?;
2786 self.conflicts = active
2787 .iter()
2788 .flat_map(|outcome| outcome.results.iter())
2789 .filter_map(|result| match result {
2790 CommitOperationOutcome::Conflict { conflict } => Some(conflict.clone()),
2791 _ => None,
2792 })
2793 .collect();
2794 self.rejections = active
2795 .iter()
2796 .flat_map(|outcome| outcome.results.iter())
2797 .filter_map(|result| match result {
2798 CommitOperationOutcome::Error { rejection } => Some(rejection.clone()),
2799 _ => None,
2800 })
2801 .collect();
2802
2803 self.lease_state = self
2804 .get_meta(LEASE_STATE_KEY)
2805 .map(|raw| serde_json::from_str(&raw))
2806 .transpose()
2807 .map_err(|error| format!("invalid persisted lease state: {error}"))?;
2808 self.schema_floor = self
2809 .get_meta(SCHEMA_FLOOR_KEY)
2810 .map(|raw| serde_json::from_str(&raw))
2811 .transpose()
2812 .map_err(|error| format!("invalid persisted schema floor: {error}"))?;
2813 self.stopped = self.schema_floor.is_some();
2814 if self.get_meta(SECURITY_PREFLIGHT_PENDING_KEY).as_deref() == Some("1") {
2817 self.security_preflight = true;
2818 }
2819 Ok(())
2820 }
2821
2822 fn prune_unknown_subscriptions(&mut self) -> Result<(), String> {
2825 let valid_tables: BTreeSet<String> = self
2826 .schema
2827 .tables
2828 .iter()
2829 .map(|table| table.name.clone())
2830 .collect();
2831 let stale_ids: Vec<String> = self
2832 .subs
2833 .iter()
2834 .filter(|sub| !valid_tables.contains(&sub.table))
2835 .map(|sub| sub.id.clone())
2836 .collect();
2837 for id in &stale_ids {
2838 self.conn
2839 .execute(
2840 "DELETE FROM _syncular_windows WHERE sub_id = ?1",
2841 rusqlite::params![id],
2842 )
2843 .map_err(|error| error.to_string())?;
2844 self.conn
2845 .execute(
2846 "DELETE FROM _syncular_window_pending_evict WHERE sub_id = ?1",
2847 rusqlite::params![id],
2848 )
2849 .map_err(|error| error.to_string())?;
2850 self.conn
2851 .execute(
2852 "DELETE FROM _syncular_subscriptions WHERE id = ?1",
2853 rusqlite::params![id],
2854 )
2855 .map_err(|error| error.to_string())?;
2856 }
2857 self.subs.retain(|sub| valid_tables.contains(&sub.table));
2858 Ok(())
2859 }
2860
2861 #[must_use]
2862 pub fn local_revision(&self) -> u64 {
2863 self.get_meta(LOCAL_REVISION_KEY)
2864 .and_then(|value| value.parse().ok())
2865 .unwrap_or(0)
2866 }
2867
2868 #[must_use]
2869 pub fn status_snapshot(&self) -> SyncStatusSnapshot {
2870 SyncStatusSnapshot {
2871 current_schema_version: self.schema.version,
2872 outbox: self.outbox.len(),
2873 upgrading: self.upgrading,
2874 lease_state: self.lease_state.clone(),
2875 schema_floor: self.schema_floor.clone(),
2876 sync_needed: self.sync_needed,
2877 }
2878 }
2879
2880 pub fn diagnostics_snapshot(
2881 &self,
2882 request: &ClientDiagnosticsRequest,
2883 ) -> Result<ClientDiagnosticsSnapshot, String> {
2884 if request.expected_subscriptions.len() > MAX_DIAGNOSTIC_EXPECTED_SUBSCRIPTIONS {
2885 return Err(format!(
2886 "sync.invalid_request: diagnosticsSnapshot accepts at most {MAX_DIAGNOSTIC_EXPECTED_SUBSCRIPTIONS} expected subscriptions"
2887 ));
2888 }
2889 let mut subscriptions = BTreeMap::<String, DiagnosticSubscription>::new();
2890 for sub in &self.subs {
2891 let reset = sub.cursor < 0 && sub.reason_code.as_deref() == Some("sync.cursor_expired");
2892 let complete =
2893 sub.state == SubState::Active && sub.cursor >= 0 && sub.bootstrap_state.is_none();
2894 let state = match sub.state {
2895 SubState::Revoked => "revoked",
2896 SubState::Failed => "failed",
2897 SubState::Active if reset => "reset",
2898 SubState::Active if complete => "complete",
2899 SubState::Active => "bootstrapping",
2900 };
2901 subscriptions.insert(
2902 sub.id.clone(),
2903 DiagnosticSubscription {
2904 id: sub.id.clone(),
2905 table: sub.table.clone(),
2906 state: state.to_owned(),
2907 complete,
2908 cursor: Some(sub.cursor),
2909 reason_code: sub.reason_code.as_deref().map(Self::diagnostic_code),
2910 },
2911 );
2912 }
2913 for expected in &request.expected_subscriptions {
2914 if expected.id.is_empty() || expected.table.is_empty() {
2915 return Err("sync.invalid_request: diagnosticsSnapshot expected subscriptions require non-empty id and table strings".to_owned());
2916 }
2917 if subscriptions
2918 .get(&expected.id)
2919 .is_some_and(|registered| registered.table != expected.table)
2920 {
2921 subscriptions.insert(
2922 expected.id.clone(),
2923 DiagnosticSubscription {
2924 id: expected.id.clone(),
2925 table: expected.table.clone(),
2926 state: "failed".to_owned(),
2927 complete: false,
2928 cursor: None,
2929 reason_code: Some("client.subscription_intent_mismatch".to_owned()),
2930 },
2931 );
2932 } else {
2933 subscriptions.entry(expected.id.clone()).or_insert_with(|| {
2934 DiagnosticSubscription {
2935 id: expected.id.clone(),
2936 table: expected.table.clone(),
2937 state: "unregistered".to_owned(),
2938 complete: false,
2939 cursor: None,
2940 reason_code: None,
2941 }
2942 });
2943 }
2944 }
2945 let mut ordered_subscriptions = Vec::new();
2946 let mut included = BTreeSet::new();
2947 for expected in &request.expected_subscriptions {
2948 if included.insert(expected.id.clone()) {
2949 if let Some(subscription) = subscriptions.get(&expected.id) {
2950 ordered_subscriptions.push(subscription.clone());
2951 }
2952 }
2953 }
2954 for (id, subscription) in &subscriptions {
2955 if included.insert(id.clone()) {
2956 ordered_subscriptions.push(subscription.clone());
2957 }
2958 }
2959 let subscriptions_truncated =
2960 ordered_subscriptions.len() > MAX_DIAGNOSTIC_EXPECTED_SUBSCRIPTIONS;
2961 ordered_subscriptions.truncate(MAX_DIAGNOSTIC_EXPECTED_SUBSCRIPTIONS);
2962 let captured_at_ms = self.clock_now_ms();
2963 let lease = if let Some(error_code) = self
2964 .lease_state
2965 .as_ref()
2966 .and_then(|state| state.error_code.clone())
2967 {
2968 ClientDiagnosticsLease {
2969 state: "stopped".to_owned(),
2970 expires_at_ms: self
2971 .lease_state
2972 .as_ref()
2973 .and_then(|state| state.expires_at_ms),
2974 error_code: Some(Self::diagnostic_code(&error_code)),
2975 }
2976 } else if let Some(expires_at_ms) = self
2977 .lease_state
2978 .as_ref()
2979 .and_then(|state| state.expires_at_ms)
2980 {
2981 ClientDiagnosticsLease {
2982 state: if expires_at_ms <= captured_at_ms {
2983 "expired".to_owned()
2984 } else {
2985 "active".to_owned()
2986 },
2987 expires_at_ms: Some(expires_at_ms),
2988 error_code: None,
2989 }
2990 } else {
2991 ClientDiagnosticsLease {
2992 state: "none".to_owned(),
2993 expires_at_ms: None,
2994 error_code: None,
2995 }
2996 };
2997 let connectivity = match self.last_round.as_ref() {
2998 Some(round) if round.status == "succeeded" => "online",
2999 Some(round)
3000 if round.status == "failed"
3001 && round
3002 .error_code
3003 .as_deref()
3004 .is_some_and(Self::retryable_transport_code) =>
3005 {
3006 "offline"
3007 }
3008 _ => "unknown",
3009 };
3010 Ok(ClientDiagnosticsSnapshot {
3011 version: CLIENT_DIAGNOSTICS_VERSION,
3012 captured_at_ms,
3013 host: ClientDiagnosticsHost {
3014 kind: "direct".to_owned(),
3015 role: "single".to_owned(),
3016 connectivity: connectivity.to_owned(),
3017 realtime: if self.realtime_connected {
3018 "connected".to_owned()
3019 } else {
3020 "disconnected".to_owned()
3021 },
3022 },
3023 security_lifecycle: self.security_lifecycle().to_owned(),
3024 schema: ClientDiagnosticsSchema {
3025 current_version: self.schema.version,
3026 upgrading: self.upgrading,
3027 required_version: self
3028 .schema_floor
3029 .as_ref()
3030 .and_then(|floor| floor.required_schema_version),
3031 latest_version: self
3032 .schema_floor
3033 .as_ref()
3034 .and_then(|floor| floor.latest_schema_version),
3035 },
3036 replica: ClientDiagnosticsReplica {
3037 local_revision: self.local_revision().to_string(),
3038 sync_needed: self.sync_needed,
3039 pending_outbox: self.outbox.len(),
3040 },
3041 lease,
3042 subscriptions: ordered_subscriptions,
3043 subscriptions_truncated,
3044 last_round: self.last_round.clone(),
3045 last_change: self.last_change.clone(),
3046 storage: self.diagnostics_storage(),
3047 })
3048 }
3049
3050 fn diagnostics_storage(&self) -> ClientDiagnosticsStorage {
3051 let read = || -> Result<ClientDiagnosticsStorage, rusqlite::Error> {
3052 let page_count: i64 = self
3053 .conn
3054 .query_row("PRAGMA page_count", [], |row| row.get(0))?;
3055 let page_size: i64 = self
3056 .conn
3057 .query_row("PRAGMA page_size", [], |row| row.get(0))?;
3058 let outbox_bytes: i64 = self.conn.query_row(
3059 "SELECT COALESCE(SUM(LENGTH(ops_json)), 0) FROM _syncular_outbox",
3060 [],
3061 |row| row.get(0),
3062 )?;
3063 let (outcome_entries, outcome_bytes): (i64, i64) = self.conn.query_row(
3064 "SELECT COUNT(*), COALESCE(SUM(LENGTH(results_json) + COALESCE(LENGTH(operations_json), 0)), 0) FROM _syncular_commit_outcomes",
3065 [],
3066 |row| Ok((row.get(0)?, row.get(1)?)),
3067 )?;
3068 let blob_bytes = if self.schema_has_blobs() {
3069 self.conn.query_row(
3070 "SELECT COALESCE(SUM(byte_length), 0) FROM _syncular_blobs",
3071 [],
3072 |row| row.get(0),
3073 )?
3074 } else {
3075 0
3076 };
3077 let pressure = self
3078 .limits
3079 .blob_cache_max_bytes
3080 .is_some_and(|limit| blob_bytes > limit);
3081 Ok(ClientDiagnosticsStorage {
3082 status: if pressure { "pressure" } else { "healthy" }.to_owned(),
3083 database_bytes_approx: Some(page_count.saturating_mul(page_size).max(0)),
3084 pending_outbox_bytes_approx: Some(outbox_bytes.max(0)),
3085 retained_outcome_bytes_approx: Some(outcome_bytes.max(0)),
3086 retained_outcome_entries: Some(outcome_entries.max(0)),
3087 blob_cache_bytes_approx: Some(blob_bytes.max(0)),
3088 pressure_reason_code: pressure.then(|| "client.blob_cache_over_limit".to_owned()),
3089 })
3090 };
3091 read().unwrap_or_else(|_| ClientDiagnosticsStorage {
3092 status: "unreadable".to_owned(),
3093 database_bytes_approx: None,
3094 pending_outbox_bytes_approx: None,
3095 retained_outcome_bytes_approx: None,
3096 retained_outcome_entries: None,
3097 blob_cache_bytes_approx: None,
3098 pressure_reason_code: None,
3099 })
3100 }
3101
3102 pub fn drain_change_batches(&mut self) -> Vec<ClientChangeBatch> {
3103 self.change_queue.drain(..).collect()
3104 }
3105
3106 pub fn drain_sync_intents(&mut self) -> Vec<SyncIntent> {
3107 self.sync_intent_queue.drain(..).collect()
3108 }
3109
3110 fn schedule_background_retry(&mut self) {
3111 self.sync_intent_queue.push_back(SyncIntent::Background {
3112 delay_ms: self.retry_delay_ms,
3113 });
3114 self.retry_delay_ms = (self.retry_delay_ms * 2).min(30_000);
3115 }
3116
3117 fn reset_background_retry(&mut self) {
3118 self.retry_delay_ms = 250;
3119 }
3120
3121 fn retryable_transport_code(code: &str) -> bool {
3122 code == "transport.failed"
3123 || code == "transport.unavailable"
3124 || code == "sync.transport_failed"
3125 }
3126
3127 fn diagnostic_code(code: &str) -> String {
3128 let valid = !code.is_empty()
3129 && code.len() <= 96
3130 && code.contains('.')
3131 && code
3132 .bytes()
3133 .next()
3134 .is_some_and(|byte| byte.is_ascii_lowercase())
3135 && code.bytes().all(|byte| {
3136 byte.is_ascii_lowercase()
3137 || byte.is_ascii_digit()
3138 || matches!(byte, b'.' | b'_' | b'-')
3139 });
3140 if valid {
3141 code.to_owned()
3142 } else {
3143 "client.unknown_failure".to_owned()
3144 }
3145 }
3146
3147 fn set_sync_needed(&mut self, value: bool, interactive: bool) {
3148 if self.sync_needed != value {
3149 if self.begin_observation("syncular_status").is_ok() {
3150 self.sync_needed = value;
3151 let batch = ChangeAccumulator {
3152 status: true,
3153 ..ChangeAccumulator::default()
3154 };
3155 if self.finish_observation("syncular_status", batch).is_err() {
3156 self.rollback_observation("syncular_status");
3157 }
3158 } else {
3159 self.sync_needed = value;
3160 }
3161 }
3162 if value && interactive {
3163 self.sync_intent_queue.push_back(SyncIntent::Interactive);
3164 }
3165 }
3166
3167 fn begin_observation(&self, name: &str) -> Result<(), String> {
3168 self.conn
3169 .execute_batch(&format!("SAVEPOINT {name}"))
3170 .map_err(|error| error.to_string())
3171 }
3172
3173 fn rollback_observation(&self, name: &str) {
3174 let _ = self
3175 .conn
3176 .execute_batch(&format!("ROLLBACK TO {name}; RELEASE {name}"));
3177 }
3178
3179 fn finish_observation(&mut self, name: &str, batch: ChangeAccumulator) -> Result<(), String> {
3180 if !batch.touched() {
3181 self.conn
3182 .execute_batch(&format!("RELEASE {name}"))
3183 .map_err(|error| error.to_string())?;
3184 return Ok(());
3185 }
3186 let revision = self
3187 .local_revision()
3188 .checked_add(1)
3189 .ok_or_else(|| "local revision exhausted u64".to_owned())?;
3190 self.conn
3191 .execute(
3192 "INSERT OR REPLACE INTO _syncular_meta(key, value) VALUES (?1, ?2)",
3193 rusqlite::params![LOCAL_REVISION_KEY, revision.to_string()],
3194 )
3195 .map_err(|error| error.to_string())?;
3196 let status = batch.status.then(|| self.status_snapshot());
3197 let event = ClientChangeBatch {
3198 revision: revision.to_string(),
3199 tables: batch
3200 .tables
3201 .into_iter()
3202 .map(|(table, scope_keys)| TableChange {
3203 table,
3204 scope_keys: scope_keys.map(|keys| keys.into_iter().collect()),
3205 })
3206 .collect(),
3207 windows: batch
3208 .windows
3209 .into_iter()
3210 .map(|((base_key, table), units)| WindowChange {
3211 base_key,
3212 table,
3213 units: units.into_iter().collect(),
3214 })
3215 .collect(),
3216 status,
3217 conflicts_changed: batch.conflicts,
3218 rejections_changed: batch.rejections,
3219 outcomes_changed: batch.outcomes,
3220 };
3221 self.conn
3222 .execute_batch(&format!("RELEASE {name}"))
3223 .map_err(|error| error.to_string())?;
3224 let mut diagnostic_tables = event
3225 .tables
3226 .iter()
3227 .map(|entry| entry.table.clone())
3228 .collect::<BTreeSet<_>>()
3229 .into_iter()
3230 .collect::<Vec<_>>();
3231 let mut diagnostic_windows = event
3232 .windows
3233 .iter()
3234 .map(|entry| entry.table.clone())
3235 .collect::<BTreeSet<_>>()
3236 .into_iter()
3237 .collect::<Vec<_>>();
3238 let domains_truncated = diagnostic_tables.len() > MAX_DIAGNOSTIC_DOMAINS
3239 || diagnostic_windows.len() > MAX_DIAGNOSTIC_DOMAINS;
3240 diagnostic_tables.truncate(MAX_DIAGNOSTIC_DOMAINS);
3241 diagnostic_windows.truncate(MAX_DIAGNOSTIC_DOMAINS);
3242 self.last_change = Some(DiagnosticLastChange {
3243 revision: event.revision.clone(),
3244 recorded_at_ms: self.clock_now_ms(),
3245 tables: diagnostic_tables,
3246 windows: diagnostic_windows,
3247 domains_truncated,
3248 status_changed: event.status.is_some(),
3249 conflicts_changed: event.conflicts_changed,
3250 rejections_changed: event.rejections_changed,
3251 outcomes_changed: event.outcomes_changed,
3252 });
3253 self.change_queue.push_back(event);
3254 Ok(())
3255 }
3256
3257 fn record_scope_map(
3258 &self,
3259 batch: &mut ChangeAccumulator,
3260 table_name: &str,
3261 scopes: &[(String, Vec<String>)],
3262 ) {
3263 let Some(table) = self.schema.table(table_name) else {
3264 return;
3265 };
3266 for (variable, values) in scopes {
3267 let Some(scope) = table
3268 .scope_variables
3269 .iter()
3270 .find(|scope| &scope.variable == variable)
3271 else {
3272 continue;
3273 };
3274 for value in values {
3275 batch.scope(table_name, format!("{}:{value}", scope.prefix));
3276 }
3277 }
3278 }
3279
3280 fn record_row_scopes(
3282 &self,
3283 batch: &mut ChangeAccumulator,
3284 table_name: &str,
3285 row_id: &str,
3286 base: bool,
3287 ) -> bool {
3288 let Some(table) = self.schema.table(table_name) else {
3289 return false;
3290 };
3291 if table.scope_variables.is_empty() {
3292 return false;
3293 }
3294 let columns = table
3295 .scope_variables
3296 .iter()
3297 .map(|scope| quote_ident(&scope.column))
3298 .collect::<Vec<_>>()
3299 .join(", ");
3300 let full_table = if base {
3301 base_table(table_name)
3302 } else {
3303 visible_table(table_name)
3304 };
3305 let sql = format!(
3306 "SELECT {columns} FROM {full_table} WHERE CAST({} AS TEXT) = ?1 LIMIT 1",
3307 quote_ident(&table.primary_key)
3308 );
3309 let Ok(mut stmt) = self.conn.prepare(&sql) else {
3310 return false;
3311 };
3312 let values = stmt.query_row(rusqlite::params![row_id], |row| {
3313 let mut values = Vec::with_capacity(table.scope_variables.len());
3314 for index in 0..table.scope_variables.len() {
3315 values.push(row.get::<_, Option<String>>(index)?);
3316 }
3317 Ok(values)
3318 });
3319 let Ok(values) = values else {
3320 return false;
3321 };
3322 let mut recorded = false;
3323 for (scope, value) in table.scope_variables.iter().zip(values) {
3324 if let Some(value) = value {
3325 batch.scope(table_name, format!("{}:{value}", scope.prefix));
3326 recorded = true;
3327 }
3328 }
3329 recorded
3330 }
3331
3332 fn record_commit_changes(
3333 &self,
3334 batch: &mut ChangeAccumulator,
3335 tables: &[String],
3336 changes: &[ssp2::model::Change],
3337 ) {
3338 for change in changes {
3339 let Some(table_name) = tables.get(change.table_index as usize) else {
3340 continue;
3341 };
3342 let mut precise = self.record_row_scopes(batch, table_name, &change.row_id, true);
3343 if let Some(table) = self.schema.table(table_name) {
3344 for (variable, value) in &change.scopes {
3345 if let Some(scope) = table
3346 .scope_variables
3347 .iter()
3348 .find(|scope| &scope.variable == variable)
3349 {
3350 batch.scope(table_name, format!("{}:{value}", scope.prefix));
3351 precise = true;
3352 }
3353 }
3354 }
3355 if !precise {
3356 batch.table(table_name);
3357 }
3358 }
3359 }
3360
3361 fn scoped_rows_exist(&self, table_name: &str, effective: &[(String, Vec<String>)]) -> bool {
3362 if effective.is_empty() {
3363 return false;
3364 }
3365 let Some(table) = self.schema.table(table_name) else {
3366 return false;
3367 };
3368 let mut clauses = Vec::new();
3369 let mut params = Vec::new();
3370 for (variable, values) in effective {
3371 let Some(column) = table.scope_column(variable) else {
3372 return false;
3373 };
3374 if values.is_empty() {
3375 return false;
3376 }
3377 let placeholders = values
3378 .iter()
3379 .map(|value| {
3380 params.push(SqlValue::Text(value.clone()));
3381 "?"
3382 })
3383 .collect::<Vec<_>>()
3384 .join(", ");
3385 clauses.push(format!("{} IN ({placeholders})", quote_ident(column)));
3386 }
3387 let sql = format!(
3388 "SELECT 1 FROM {} WHERE {} LIMIT 1",
3389 base_table(table_name),
3390 clauses.join(" AND ")
3391 );
3392 self.conn
3393 .query_row(&sql, rusqlite::params_from_iter(params), |_| Ok(()))
3394 .is_ok()
3395 }
3396
3397 pub fn upgrading(&self) -> bool {
3399 self.upgrading
3400 }
3401
3402 pub fn recreate_with_schema(&mut self, schema_json: &Value) -> Result<(), String> {
3408 let new_schema = parse_schema_json(schema_json)?;
3409 let marker: Option<i32> = self
3410 .get_meta(LOCAL_SCHEMA_VERSION_KEY)
3411 .and_then(|v| v.parse().ok());
3412 self.schema = new_schema;
3413 if marker != Some(self.schema.version) {
3414 self.run_schema_reset()?;
3415 }
3416 self.prune_unknown_subscriptions()?;
3417 self.enqueue_startup_sync_if_needed();
3421 Ok(())
3422 }
3423
3424 fn enqueue_startup_sync_if_needed(&mut self) {
3425 let startup_work = !self.stopped
3426 && (!self.outbox.is_empty()
3427 || self.subs.iter().any(|sub| sub.state == SubState::Active));
3428 if startup_work {
3429 self.sync_needed = true;
3430 self.sync_intent_queue.push_back(SyncIntent::Interactive);
3431 }
3432 }
3433
3434 fn clear_satisfied_persisted_schema_floor(&mut self) {
3441 let satisfied = self
3442 .schema_floor
3443 .as_ref()
3444 .and_then(|floor| floor.required_schema_version)
3445 .is_some_and(|required| self.schema.version >= required);
3446 if !satisfied {
3447 return;
3448 }
3449 self.schema_floor = None;
3450 self.stopped = false;
3451 self.delete_meta(SCHEMA_FLOOR_KEY);
3452 }
3453
3454 fn run_schema_reset(&mut self) -> Result<(), String> {
3460 self.begin_observation("syncular_schema_reset")?;
3461 let mut batch = ChangeAccumulator::default();
3462 let result = self.run_schema_reset_observed(&mut batch, true);
3463 if let Err(error) = result {
3464 self.rollback_observation("syncular_schema_reset");
3465 return Err(error);
3466 }
3467 if let Err(error) = self.finish_observation("syncular_schema_reset", batch) {
3468 self.rollback_observation("syncular_schema_reset");
3469 return Err(error);
3470 }
3471 Ok(())
3472 }
3473
3474 fn run_schema_reset_observed(
3475 &mut self,
3476 batch: &mut ChangeAccumulator,
3477 drop_incompatible: bool,
3478 ) -> Result<(), String> {
3479 self.upgrading = true;
3480 batch.status = true;
3481 for table in &self.schema.tables {
3482 batch.table(&table.name);
3483 }
3484 for (base_key, unit, table) in self.load_registered_window_units() {
3485 batch.window(&base_key, &table, &unit);
3486 }
3487 self.insert_sql.borrow_mut().clear();
3489 self.overlay_dirty.set(true);
3490 let virtual_tables: Vec<String> = {
3498 let mut stmt = self
3499 .conn
3500 .prepare(
3501 "SELECT name FROM sqlite_master WHERE type = 'table' AND sql LIKE 'CREATE VIRTUAL TABLE%'",
3502 )
3503 .map_err(|e| e.to_string())?;
3504 let rows = stmt
3505 .query_map([], |row| row.get::<_, String>(0))
3506 .map_err(|e| e.to_string())?;
3507 rows.filter_map(Result::ok)
3508 .filter(|name| is_synced_table_name(name))
3509 .collect()
3510 };
3511 for name in virtual_tables {
3512 self.conn
3513 .execute(&format!("DROP TABLE IF EXISTS {}", quote_ident(&name)), [])
3514 .map_err(|e| e.to_string())?;
3515 }
3516 let existing: Vec<String> = {
3517 let mut stmt = self
3518 .conn
3519 .prepare("SELECT name FROM sqlite_master WHERE type = 'table'")
3520 .map_err(|e| e.to_string())?;
3521 let rows = stmt
3522 .query_map([], |row| row.get::<_, String>(0))
3523 .map_err(|e| e.to_string())?;
3524 rows.filter_map(Result::ok)
3525 .filter(|name| is_synced_table_name(name))
3526 .collect()
3527 };
3528 for name in &existing {
3529 if let Some(table) = name.strip_prefix("_syncular_base_") {
3530 batch.table(table);
3531 } else if !name.starts_with("_syncular_") {
3532 batch.table(name);
3533 }
3534 }
3535 for name in existing {
3536 self.conn
3537 .execute(&format!("DROP TABLE IF EXISTS {}", quote_ident(&name)), [])
3538 .map_err(|e| e.to_string())?;
3539 }
3540 self.create_synced_tables()?;
3542 for sub in &mut self.subs {
3544 sub.cursor = -1;
3545 sub.bootstrap_state = None;
3546 sub.effective = None;
3547 sub.state = SubState::Active;
3548 sub.reason_code = None;
3549 sub.synced_once = false;
3550 }
3551 let subs = self.subs.clone();
3552 for sub in &subs {
3553 self.persist_sub(sub);
3554 }
3555 self.stopped = false;
3557 self.schema_floor = None;
3558 self.delete_meta(SCHEMA_FLOOR_KEY);
3559 self.set_meta(LOCAL_SCHEMA_VERSION_KEY, &self.schema.version.to_string());
3561 if drop_incompatible && self.drop_incompatible_outbox()? {
3565 batch.rejections = true;
3566 batch.status = true;
3567 batch.outcomes = true;
3568 }
3569 self.rebuild_overlay();
3571 Ok(())
3572 }
3573
3574 fn drop_incompatible_outbox(&mut self) -> Result<bool, String> {
3578 let schema = &self.schema;
3579 let incompatible = self
3580 .outbox
3581 .iter()
3582 .filter(|commit| {
3583 commit.ops.iter().any(|op| {
3584 if !op.upsert {
3585 return false;
3586 }
3587 match schema.table(&op.table) {
3588 None => true,
3589 Some(table) => op.values.as_ref().is_some_and(|values| {
3590 values
3591 .keys()
3592 .any(|key| !table.columns.iter().any(|c| &c.name == key))
3593 }),
3594 }
3595 })
3596 })
3597 .cloned()
3598 .collect::<Vec<_>>();
3599 if incompatible.is_empty() {
3600 return Ok(false);
3601 }
3602 let mut rejections = Vec::new();
3603 for commit in &incompatible {
3604 let results = commit
3605 .ops
3606 .iter()
3607 .enumerate()
3608 .map(|(op_index, operation)| {
3609 let rejection = RejectionRecord {
3610 client_commit_id: commit.client_commit_id.clone(),
3611 op_index: op_index as i32,
3612 code: OUTBOX_INCOMPATIBLE_CODE.to_owned(),
3613 message: "the persisted commit cannot encode under the current schema"
3614 .to_owned(),
3615 retryable: false,
3616 details: None,
3617 operation: Some(CommitOperation::from(operation)),
3618 };
3619 rejections.push(rejection.clone());
3620 CommitOperationOutcome::Error { rejection }
3621 })
3622 .collect::<Vec<_>>();
3623 self.persist_commit_outcome(
3624 &commit.client_commit_id,
3625 CommitOutcomeStatus::Rejected,
3626 &results,
3627 Some(&commit.ops),
3628 )?;
3629 self.delete_outbox_persisted(&commit.client_commit_id)?;
3630 }
3631 self.prune_commit_outcomes()?;
3632 let incompatible_ids = incompatible
3633 .iter()
3634 .map(|commit| commit.client_commit_id.as_str())
3635 .collect::<BTreeSet<_>>();
3636 self.outbox
3637 .retain(|commit| !incompatible_ids.contains(commit.client_commit_id.as_str()));
3638 self.rejections.extend(rejections);
3639 Ok(true)
3640 }
3641
3642 fn create_synced_tables(&self) -> Result<(), String> {
3645 for table in &self.schema.tables {
3646 for (full, index_prefix) in [
3650 (base_table(&table.name), "_syncular_base_"),
3651 (visible_table(&table.name), ""),
3652 ] {
3653 let mut cols: Vec<String> =
3654 table.columns.iter().map(|c| quote_ident(&c.name)).collect();
3655 cols.push("\"_syncular_version\" INTEGER NOT NULL".to_owned());
3656 let sql = format!(
3657 "CREATE TABLE IF NOT EXISTS {full} ({} , PRIMARY KEY ({}))",
3658 cols.join(", "),
3659 quote_ident(&table.primary_key)
3660 );
3661 self.conn.execute(&sql, []).map_err(|e| e.to_string())?;
3662 for index in &table.indexes {
3666 let unique = if index.unique { "UNIQUE " } else { "" };
3667 let index_name = quote_ident(&format!("{index_prefix}{}", index.name));
3668 let cols_sql = index
3669 .columns
3670 .iter()
3671 .map(|c| quote_ident(c))
3672 .collect::<Vec<_>>()
3673 .join(", ");
3674 let index_sql = format!(
3675 "CREATE {unique}INDEX IF NOT EXISTS {index_name} ON {full} ({cols_sql})"
3676 );
3677 self.conn
3678 .execute(&index_sql, [])
3679 .map_err(|e| e.to_string())?;
3680 }
3681 }
3682 }
3683 for table in &self.schema.tables {
3686 for index in &table.fts_indexes {
3687 self.create_fts_projection(table, index)?;
3688 }
3689 }
3690 Ok(())
3691 }
3692
3693 fn fts_projection_exists(&self, index: &FtsIndexSchema) -> Result<bool, String> {
3694 let count: i64 = self
3695 .conn
3696 .query_row(
3697 "SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name=?1",
3698 rusqlite::params![index.name],
3699 |row| row.get(0),
3700 )
3701 .map_err(|e| e.to_string())?;
3702 Ok(count > 0)
3703 }
3704
3705 fn drop_fts_triggers(&self, index: &FtsIndexSchema) -> Result<(), String> {
3706 for suffix in ["ai", "ad", "au"] {
3707 self.conn
3708 .execute(
3709 &format!(
3710 "DROP TRIGGER IF EXISTS {}",
3711 quote_ident(&format!("{}_{suffix}", index.name))
3712 ),
3713 [],
3714 )
3715 .map_err(|e| e.to_string())?;
3716 }
3717 Ok(())
3718 }
3719
3720 fn create_fts_triggers(
3721 &self,
3722 table: &TableSchema,
3723 index: &FtsIndexSchema,
3724 ) -> Result<(), String> {
3725 let fts = quote_ident(&index.name);
3726 let source = visible_table(&table.name);
3727 let source_id = quote_ident(FTS_SOURCE_ID_COLUMN);
3728 let pk = quote_ident(&table.primary_key);
3729 let projection_columns = std::iter::once(source_id.clone())
3730 .chain(index.columns.iter().map(|column| quote_ident(column)))
3731 .collect::<Vec<_>>()
3732 .join(", ");
3733 let new_values = std::iter::once(format!("CAST(new.{pk} AS TEXT)"))
3734 .chain(
3735 index
3736 .columns
3737 .iter()
3738 .map(|column| format!("new.{}", quote_ident(column))),
3739 )
3740 .collect::<Vec<_>>()
3741 .join(", ");
3742 let delete_new = format!("DELETE FROM {fts} WHERE {source_id} = CAST(new.{pk} AS TEXT)");
3743 let delete_old = format!("DELETE FROM {fts} WHERE {source_id} = CAST(old.{pk} AS TEXT)");
3744 let insert_new = format!("INSERT INTO {fts} ({projection_columns}) VALUES ({new_values})");
3745 let bi = quote_ident(&format!("{}_bi", index.name));
3746 let ai = quote_ident(&format!("{}_ai", index.name));
3747 let ad = quote_ident(&format!("{}_ad", index.name));
3748 let au = quote_ident(&format!("{}_au", index.name));
3749 let replacement_exists = format!("EXISTS (SELECT 1 FROM {source} WHERE {pk} = new.{pk})");
3750 let sql = format!(
3751 "DROP TRIGGER IF EXISTS {bi};
3752 DROP TRIGGER IF EXISTS {ai};
3753 DROP TRIGGER IF EXISTS {ad};
3754 DROP TRIGGER IF EXISTS {au};
3755 CREATE TRIGGER {bi} BEFORE INSERT ON {source} WHEN {replacement_exists} BEGIN {delete_new}; END;
3756 CREATE TRIGGER {ai} AFTER INSERT ON {source} BEGIN {insert_new}; END;
3757 CREATE TRIGGER {ad} AFTER DELETE ON {source} BEGIN {delete_old}; END;
3758 CREATE TRIGGER {au} AFTER UPDATE ON {source} BEGIN {delete_old}; {delete_new}; {insert_new}; END;",
3759 );
3760 self.conn.execute_batch(&sql).map_err(|e| e.to_string())
3761 }
3762
3763 fn rebuild_fts_projection(
3764 &self,
3765 table: &TableSchema,
3766 index: &FtsIndexSchema,
3767 ) -> Result<(), String> {
3768 let fts = quote_ident(&index.name);
3769 let source_id = quote_ident(FTS_SOURCE_ID_COLUMN);
3770 let pk = quote_ident(&table.primary_key);
3771 let indexed_columns = index
3772 .columns
3773 .iter()
3774 .map(|column| quote_ident(column))
3775 .collect::<Vec<_>>();
3776 let projection_columns = std::iter::once(source_id)
3777 .chain(indexed_columns.iter().cloned())
3778 .collect::<Vec<_>>()
3779 .join(", ");
3780 let sql = format!(
3781 "DELETE FROM {fts}; INSERT INTO {fts} ({projection_columns}) SELECT CAST({pk} AS TEXT), {columns} FROM {source};",
3782 columns = indexed_columns.join(", "),
3783 source = visible_table(&table.name),
3784 );
3785 self.conn.execute_batch(&sql).map_err(|e| e.to_string())
3786 }
3787
3788 fn create_fts_projection(
3789 &self,
3790 table: &TableSchema,
3791 index: &FtsIndexSchema,
3792 ) -> Result<(), String> {
3793 let existed = self.fts_projection_exists(index)?;
3794 let tokenizer = index.tokenize.replace('\'', "''");
3795 let columns = index
3796 .columns
3797 .iter()
3798 .map(|column| quote_ident(column))
3799 .collect::<Vec<_>>()
3800 .join(", ");
3801 let sql = format!(
3802 "CREATE VIRTUAL TABLE IF NOT EXISTS {fts} USING fts5({source_id} UNINDEXED, {columns}, tokenize='{tokenizer}')",
3803 fts = quote_ident(&index.name),
3804 source_id = quote_ident(FTS_SOURCE_ID_COLUMN),
3805 );
3806 self.conn.execute(&sql, []).map_err(|error| {
3807 format!(
3808 "cannot create local FTS5 projection {:?}: {error}",
3809 index.name
3810 )
3811 })?;
3812 self.create_fts_triggers(table, index)?;
3813 if !existed {
3814 self.rebuild_fts_projection(table, index)?;
3815 }
3816 Ok(())
3817 }
3818
3819 fn persist_sub(&self, sub: &Subscription) {
3822 let state = serde_json::json!({
3823 "requested": scope_map_to_json(&sub.requested),
3824 "params": sub.params,
3825 "cursor": sub.cursor,
3826 "bootstrapState": sub.bootstrap_state,
3827 "status": sub.state.name(),
3828 "reasonCode": sub.reason_code,
3829 "effectiveScopes": sub.effective.as_ref().map(|e| scope_map_to_json(e)),
3830 "syncedOnce": sub.synced_once,
3831 });
3832 let _ = self.conn.execute(
3833 "INSERT OR REPLACE INTO _syncular_subscriptions (id, tbl, state_json) VALUES (?1, ?2, ?3)",
3834 rusqlite::params![sub.id, sub.table, state.to_string()],
3835 );
3836 }
3837
3838 fn persist_outbox_insert(&self, commit: &OutboxCommit) {
3839 let ops: Vec<Value> = commit
3840 .ops
3841 .iter()
3842 .map(|op| {
3843 serde_json::json!({
3844 "op": if op.upsert { "upsert" } else { "delete" },
3845 "table": op.table,
3846 "rowId": op.row_id,
3847 "baseVersion": op.base_version,
3848 "values": op.values.clone().map(Value::Object),
3849 "changedFields": op.changed_fields,
3850 })
3851 })
3852 .collect();
3853 let _ = self.conn.execute(
3854 "INSERT OR REPLACE INTO _syncular_outbox (commit_id, ops_json) VALUES (?1, ?2)",
3855 rusqlite::params![commit.client_commit_id, Value::Array(ops).to_string()],
3856 );
3857 }
3858
3859 fn delete_outbox_persisted(&self, client_commit_id: &str) -> Result<(), String> {
3860 self.conn
3861 .execute(
3862 "DELETE FROM _syncular_outbox WHERE commit_id = ?1",
3863 rusqlite::params![client_commit_id],
3864 )
3865 .map(|_| ())
3866 .map_err(|error| error.to_string())
3867 }
3868
3869 fn outcome_status_name(status: CommitOutcomeStatus) -> &'static str {
3870 match status {
3871 CommitOutcomeStatus::Applied => "applied",
3872 CommitOutcomeStatus::Cached => "cached",
3873 CommitOutcomeStatus::Conflict => "conflict",
3874 CommitOutcomeStatus::Rejected => "rejected",
3875 }
3876 }
3877
3878 fn outcome_resolution_name(resolution: CommitOutcomeResolution) -> &'static str {
3879 match resolution {
3880 CommitOutcomeResolution::Active => "active",
3881 CommitOutcomeResolution::ResolvedKeepServer => "resolved_keep_server",
3882 CommitOutcomeResolution::Superseded => "superseded",
3883 CommitOutcomeResolution::Dismissed => "dismissed",
3884 }
3885 }
3886
3887 fn parse_outcome_status(value: &str) -> Result<CommitOutcomeStatus, String> {
3888 match value {
3889 "applied" => Ok(CommitOutcomeStatus::Applied),
3890 "cached" => Ok(CommitOutcomeStatus::Cached),
3891 "conflict" => Ok(CommitOutcomeStatus::Conflict),
3892 "rejected" => Ok(CommitOutcomeStatus::Rejected),
3893 _ => Err(format!("invalid persisted commit outcome status {value:?}")),
3894 }
3895 }
3896
3897 fn parse_outcome_resolution(value: &str) -> Result<CommitOutcomeResolution, String> {
3898 match value {
3899 "active" => Ok(CommitOutcomeResolution::Active),
3900 "resolved_keep_server" => Ok(CommitOutcomeResolution::ResolvedKeepServer),
3901 "superseded" => Ok(CommitOutcomeResolution::Superseded),
3902 "dismissed" => Ok(CommitOutcomeResolution::Dismissed),
3903 _ => Err(format!(
3904 "invalid persisted commit outcome resolution {value:?}"
3905 )),
3906 }
3907 }
3908
3909 fn persist_commit_outcome(
3910 &self,
3911 client_commit_id: &str,
3912 status: CommitOutcomeStatus,
3913 results: &[CommitOperationOutcome],
3914 operations: Option<&[OutboxOp]>,
3915 ) -> Result<(), String> {
3916 let results_json = serde_json::to_string(results).map_err(|error| error.to_string())?;
3917 let operations_json = operations
3918 .map(|items| {
3919 serde_json::to_string(&items.iter().map(CommitOperation::from).collect::<Vec<_>>())
3920 })
3921 .transpose()
3922 .map_err(|error| error.to_string())?;
3923 self.conn
3924 .execute(
3925 "INSERT INTO _syncular_commit_outcomes (
3926 client_commit_id, status, recorded_at_ms, results_json,
3927 operations_json, resolution
3928 ) VALUES (?1, ?2, ?3, ?4, ?5, 'active')",
3929 rusqlite::params![
3930 client_commit_id,
3931 Self::outcome_status_name(status),
3932 self.clock_now_ms(),
3933 results_json,
3934 operations_json
3935 ],
3936 )
3937 .map(|_| ())
3938 .map_err(|error| error.to_string())
3939 }
3940
3941 fn outcome_from_row(row: StoredCommitOutcomeRow) -> Result<CommitOutcome, String> {
3942 let StoredCommitOutcomeRow {
3943 sequence,
3944 client_commit_id,
3945 status,
3946 recorded_at_ms,
3947 results_json,
3948 operations_json,
3949 resolution,
3950 resolved_at_ms,
3951 replacement_client_commit_id,
3952 } = row;
3953 Ok(CommitOutcome {
3954 sequence,
3955 client_commit_id,
3956 status: Self::parse_outcome_status(&status)?,
3957 recorded_at_ms,
3958 results: serde_json::from_str(&results_json)
3959 .map_err(|error| format!("invalid persisted commit outcome results: {error}"))?,
3960 operations: operations_json
3961 .map(|value| {
3962 serde_json::from_str(&value).map_err(|error| {
3963 format!("invalid persisted commit outcome operations: {error}")
3964 })
3965 })
3966 .transpose()?,
3967 resolution: Self::parse_outcome_resolution(&resolution)?,
3968 resolved_at_ms,
3969 replacement_client_commit_id,
3970 })
3971 }
3972
3973 pub fn commit_outcome(&self, client_commit_id: &str) -> Result<Option<CommitOutcome>, String> {
3974 let row = self
3975 .conn
3976 .query_row(
3977 "SELECT seq, client_commit_id, status, recorded_at_ms, results_json, operations_json,
3978 resolution, resolved_at_ms, replacement_client_commit_id
3979 FROM _syncular_commit_outcomes WHERE client_commit_id = ?1",
3980 rusqlite::params![client_commit_id],
3981 |row| {
3982 Ok(StoredCommitOutcomeRow {
3983 sequence: row.get(0)?,
3984 client_commit_id: row.get(1)?,
3985 status: row.get(2)?,
3986 recorded_at_ms: row.get(3)?,
3987 results_json: row.get(4)?,
3988 operations_json: row.get(5)?,
3989 resolution: row.get(6)?,
3990 resolved_at_ms: row.get(7)?,
3991 replacement_client_commit_id: row.get(8)?,
3992 })
3993 },
3994 )
3995 .optional()
3996 .map_err(|error| error.to_string())?;
3997 row.map(Self::outcome_from_row).transpose()
3998 }
3999
4000 pub fn commit_outcomes(&self, query: CommitOutcomeQuery) -> Result<Vec<CommitOutcome>, String> {
4001 if query.limit == Some(0) {
4002 return Err("sync.invalid_request: commit outcome limit must be positive".to_owned());
4003 }
4004 let mut sql = String::from(
4005 "SELECT seq, client_commit_id, status, recorded_at_ms, results_json, operations_json,
4006 resolution, resolved_at_ms, replacement_client_commit_id
4007 FROM _syncular_commit_outcomes",
4008 );
4009 if query.active_only {
4010 sql.push_str(" WHERE resolution = 'active' AND status IN ('conflict', 'rejected')");
4011 }
4012 sql.push_str(" ORDER BY seq DESC");
4013 if let Some(limit) = query.limit {
4014 sql.push_str(&format!(" LIMIT {limit}"));
4015 }
4016 let mut stmt = self.conn.prepare(&sql).map_err(|error| error.to_string())?;
4017 let rows = stmt
4018 .query_map([], |row| {
4019 Ok(StoredCommitOutcomeRow {
4020 sequence: row.get(0)?,
4021 client_commit_id: row.get(1)?,
4022 status: row.get(2)?,
4023 recorded_at_ms: row.get(3)?,
4024 results_json: row.get(4)?,
4025 operations_json: row.get(5)?,
4026 resolution: row.get(6)?,
4027 resolved_at_ms: row.get(7)?,
4028 replacement_client_commit_id: row.get(8)?,
4029 })
4030 })
4031 .map_err(|error| error.to_string())?;
4032 let mut outcomes = Vec::new();
4033 for row in rows {
4034 outcomes.push(Self::outcome_from_row(
4035 row.map_err(|error| error.to_string())?,
4036 )?);
4037 }
4038 Ok(outcomes)
4039 }
4040
4041 fn prune_commit_outcomes(&self) -> Result<(), String> {
4042 #[cfg(test)]
4043 self.outcome_prune_count
4044 .set(self.outcome_prune_count.get() + 1);
4045 let max_entries = self.limits.outcome_retention_max_entries.unwrap_or(1_000);
4046 let count = self
4047 .conn
4048 .query_row(
4049 "SELECT COUNT(*) FROM _syncular_commit_outcomes",
4050 [],
4051 |row| row.get::<_, i64>(0),
4052 )
4053 .map_err(|error| error.to_string())? as usize;
4054 let excess = count.saturating_sub(max_entries);
4055 if excess == 0 {
4056 return Ok(());
4057 }
4058 let mut stmt = self
4059 .conn
4060 .prepare(
4061 "SELECT seq FROM _syncular_commit_outcomes
4062 WHERE status IN ('applied', 'cached') OR resolution != 'active'
4063 ORDER BY seq ASC LIMIT ?1",
4064 )
4065 .map_err(|error| error.to_string())?;
4066 let rows = stmt
4067 .query_map(rusqlite::params![excess as i64], |row| row.get::<_, i64>(0))
4068 .map_err(|error| error.to_string())?;
4069 let sequences = rows
4070 .collect::<Result<Vec<_>, _>>()
4071 .map_err(|error| error.to_string())?;
4072 drop(stmt);
4073 for sequence in sequences {
4074 self.conn
4075 .execute(
4076 "DELETE FROM _syncular_commit_outcomes WHERE seq = ?1",
4077 rusqlite::params![sequence],
4078 )
4079 .map_err(|error| error.to_string())?;
4080 }
4081 Ok(())
4082 }
4083
4084 pub fn subscribe(
4087 &mut self,
4088 id: String,
4089 table: String,
4090 scopes: Vec<(String, Vec<String>)>,
4091 params: Option<String>,
4092 ) -> Result<(), String> {
4093 if self.schema.table(&table).is_none() {
4094 return Err(format!("unknown table {table:?}"));
4095 }
4096 if let Some(existing) = self.subs.iter().find(|subscription| subscription.id == id) {
4097 let same_intent = existing.table == table
4098 && canonical_scope_json(&existing.requested) == canonical_scope_json(&scopes)
4099 && existing.params == params;
4100 if same_intent {
4101 return Ok(());
4102 }
4103 return Err(
4104 "client.subscription_intent_mismatch: the subscription id is already registered for a different table, scopes, or params"
4105 .to_owned(),
4106 );
4107 }
4108 let sub = Subscription {
4109 id: id.clone(),
4110 table,
4111 requested: scopes,
4112 params,
4113 cursor: -1,
4114 bootstrap_state: None,
4115 state: SubState::Active,
4116 reason_code: None,
4117 effective: None,
4118 synced_once: false,
4119 };
4120 self.persist_sub(&sub);
4121 self.subs.push(sub);
4122 Ok(())
4123 }
4124
4125 pub fn unsubscribe(&mut self, id: &str) {
4126 self.subs.retain(|s| s.id != id);
4127 let _ = self.conn.execute(
4128 "DELETE FROM _syncular_subscriptions WHERE id = ?1",
4129 rusqlite::params![id],
4130 );
4131 }
4132
4133 pub fn set_window(
4141 &mut self,
4142 base: &WindowBase,
4143 units: &[String],
4144 ) -> Result<CommandEffects, String> {
4145 let table = self
4146 .schema
4147 .table(&base.table)
4148 .ok_or_else(|| format!("unknown table {:?}", base.table))?;
4149 if table.scope_column(&base.variable).is_none() {
4150 return Err(format!(
4151 "setWindow: table {:?} has no scope variable {:?} (§4.8)",
4152 base.table, base.variable
4153 ));
4154 }
4155 let base_key = window_base_key(base);
4156 let wanted: std::collections::HashSet<&String> = units.iter().collect();
4157 let live = self.load_window_units(&base_key);
4158 self.begin_observation("syncular_window")?;
4159 let mut batch = ChangeAccumulator::default();
4160 let mut changed = false;
4161
4162 for unit in units {
4164 if live.iter().any(|(u, _)| u == unit) {
4165 continue;
4166 }
4167 let sub_id = derive_sub_id(base, unit);
4168 self.delete_pending_evict(&sub_id);
4169 self.insert_window_unit(&base_key, unit, &sub_id);
4170 self.subscribe(
4171 sub_id,
4172 base.table.clone(),
4173 unit_scopes(base, unit),
4174 base.params.clone(),
4175 )?;
4176 batch.window(&base_key, &base.table, unit);
4177 changed = true;
4178 }
4179
4180 for (unit, sub_id) in live {
4182 if wanted.contains(&unit) {
4183 continue;
4184 }
4185 let effective = self
4186 .subs
4187 .iter()
4188 .find(|sub| sub.id == sub_id)
4189 .and_then(|sub| sub.effective.clone())
4190 .unwrap_or_else(|| unit_scopes(base, &unit));
4191 self.record_scope_map(&mut batch, &base.table, &effective);
4192 batch.window(&base_key, &base.table, &unit);
4193 self.evict_unit(&base_key, base, &unit, &sub_id);
4194 changed = true;
4195 }
4196 if let Err(error) = self.finish_observation("syncular_window", batch) {
4197 self.rollback_observation("syncular_window");
4198 return Err(error);
4199 }
4200 Ok(if changed {
4201 CommandEffects::interactive()
4202 } else {
4203 CommandEffects::none()
4204 })
4205 }
4206
4207 pub fn window_state(&self, base: &WindowBase) -> WindowState {
4212 let mut units = Vec::new();
4213 let mut pending = Vec::new();
4214 for (unit, sub_id) in self.load_window_units(&window_base_key(base)) {
4215 let is_pending = match self.subs.iter().find(|s| s.id == sub_id) {
4216 Some(sub) => {
4217 sub.state != SubState::Active || sub.cursor < 0 || sub.bootstrap_state.is_some()
4218 }
4219 None => true,
4220 };
4221 if is_pending {
4222 pending.push(unit.clone());
4223 }
4224 units.push(unit);
4225 }
4226 WindowState { units, pending }
4227 }
4228
4229 fn load_window_units(&self, base_key: &str) -> Vec<(String, String)> {
4230 let mut stmt = match self
4231 .conn
4232 .prepare("SELECT unit, sub_id FROM _syncular_windows WHERE base = ?1 ORDER BY unit ASC")
4233 {
4234 Ok(stmt) => stmt,
4235 Err(_) => return Vec::new(),
4236 };
4237 let rows = stmt.query_map(rusqlite::params![base_key], |row| {
4238 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
4239 });
4240 match rows {
4241 Ok(rows) => rows.filter_map(Result::ok).collect(),
4242 Err(_) => Vec::new(),
4243 }
4244 }
4245
4246 fn load_registered_window_units(&self) -> Vec<(String, String, String)> {
4247 let mut stmt = match self.conn.prepare(
4248 "SELECT windows.base, windows.unit, subscriptions.tbl
4249 FROM _syncular_windows AS windows
4250 JOIN _syncular_subscriptions AS subscriptions
4251 ON subscriptions.id = windows.sub_id
4252 ORDER BY windows.base, windows.unit",
4253 ) {
4254 Ok(stmt) => stmt,
4255 Err(_) => return Vec::new(),
4256 };
4257 let rows = stmt.query_map([], |row| {
4258 Ok((
4259 row.get::<_, String>(0)?,
4260 row.get::<_, String>(1)?,
4261 row.get::<_, String>(2)?,
4262 ))
4263 });
4264 match rows {
4265 Ok(rows) => rows.filter_map(Result::ok).collect(),
4266 Err(_) => Vec::new(),
4267 }
4268 }
4269
4270 fn window_unit_by_sub_id(&self, sub_id: &str) -> Option<(String, String)> {
4271 self.conn
4272 .query_row(
4273 "SELECT base, unit FROM _syncular_windows WHERE sub_id = ?1 LIMIT 1",
4274 rusqlite::params![sub_id],
4275 |row| Ok((row.get(0)?, row.get(1)?)),
4276 )
4277 .ok()
4278 }
4279
4280 fn insert_window_unit(&self, base_key: &str, unit: &str, sub_id: &str) {
4281 let _ = self.conn.execute(
4282 "INSERT OR REPLACE INTO _syncular_windows(base, unit, sub_id) VALUES (?1, ?2, ?3)",
4283 rusqlite::params![base_key, unit, sub_id],
4284 );
4285 }
4286
4287 fn delete_window_unit(&self, base_key: &str, unit: &str) {
4288 let _ = self.conn.execute(
4289 "DELETE FROM _syncular_windows WHERE base = ?1 AND unit = ?2",
4290 rusqlite::params![base_key, unit],
4291 );
4292 }
4293
4294 fn evict_unit(&mut self, base_key: &str, base: &WindowBase, unit: &str, sub_id: &str) {
4301 let effective = self
4302 .subs
4303 .iter()
4304 .find(|s| s.id == sub_id)
4305 .and_then(|s| s.effective.clone())
4306 .unwrap_or_else(|| unit_scopes(base, unit));
4307 let pinned = self.pinned_row_ids(&base.table);
4308 let deferred = self
4309 .evict_scope_rows(&base.table, &effective, &pinned)
4310 .unwrap_or(false);
4311 self.delete_window_unit(base_key, unit);
4312 self.unsubscribe(sub_id);
4313 if deferred {
4314 self.save_pending_evict(sub_id, &base.table, &effective);
4315 } else {
4316 self.delete_pending_evict(sub_id);
4317 }
4318 self.rebuild_overlay();
4319 }
4320
4321 fn evict_scope_rows(
4325 &mut self,
4326 table_name: &str,
4327 effective: &[(String, Vec<String>)],
4328 pinned: &std::collections::HashSet<String>,
4329 ) -> Result<bool, ()> {
4330 if effective.is_empty() {
4331 return Ok(false);
4332 }
4333 let table = self.schema.table(table_name).ok_or(())?.clone();
4334 let mut clauses = Vec::new();
4335 let mut params: Vec<SqlValue> = Vec::new();
4336 for (variable, values) in effective {
4337 let column = table.scope_column(variable).ok_or(())?;
4338 let placeholders: Vec<String> = values
4339 .iter()
4340 .map(|v| {
4341 params.push(SqlValue::Text(v.clone()));
4342 "?".to_owned()
4343 })
4344 .collect();
4345 clauses.push(format!(
4346 "{} IN ({})",
4347 quote_ident(column),
4348 placeholders.join(", ")
4349 ));
4350 }
4351 let mut sql = format!(
4352 "DELETE FROM {} WHERE {}",
4353 base_table(table_name),
4354 clauses.join(" AND ")
4355 );
4356 if !pinned.is_empty() {
4357 let pk = quote_ident(&table.primary_key);
4358 let holes: Vec<String> = pinned
4359 .iter()
4360 .map(|id| {
4361 params.push(SqlValue::Text(id.clone()));
4362 "?".to_owned()
4363 })
4364 .collect();
4365 sql.push_str(&format!(" AND {} NOT IN ({})", pk, holes.join(", ")));
4366 }
4367 self.overlay_dirty.set(true);
4368 self.conn
4369 .execute(&sql, rusqlite::params_from_iter(params))
4370 .map_err(|_| ())?;
4371 if pinned.is_empty() {
4372 return Ok(false);
4373 }
4374 let mut where_params: Vec<SqlValue> = Vec::new();
4377 let mut where_clauses = Vec::new();
4378 for (variable, values) in effective {
4379 let column = table.scope_column(variable).ok_or(())?;
4380 let placeholders: Vec<String> = values
4381 .iter()
4382 .map(|v| {
4383 where_params.push(SqlValue::Text(v.clone()));
4384 "?".to_owned()
4385 })
4386 .collect();
4387 where_clauses.push(format!(
4388 "{} IN ({})",
4389 quote_ident(column),
4390 placeholders.join(", ")
4391 ));
4392 }
4393 let pk = quote_ident(&table.primary_key);
4394 let select = format!(
4395 "SELECT {} FROM {} WHERE {}",
4396 pk,
4397 base_table(table_name),
4398 where_clauses.join(" AND ")
4399 );
4400 let mut stmt = self.conn.prepare(&select).map_err(|_| ())?;
4401 let survivors: Vec<String> = stmt
4402 .query_map(rusqlite::params_from_iter(where_params), |row| {
4403 row.get::<_, String>(0)
4404 })
4405 .map_err(|_| ())?
4406 .filter_map(Result::ok)
4407 .collect();
4408 Ok(survivors.iter().any(|id| pinned.contains(id)))
4409 }
4410
4411 fn drain_pending_evictions(&mut self) {
4415 let pending = self.load_pending_evictions();
4416 if pending.is_empty() {
4417 return;
4418 }
4419 for (sub_id, table_name, effective) in pending {
4420 if self.schema.table(&table_name).is_none() {
4421 self.delete_pending_evict(&sub_id);
4422 continue;
4423 }
4424 let pinned = self.pinned_row_ids(&table_name);
4425 let deferred = self
4426 .evict_scope_rows(&table_name, &effective, &pinned)
4427 .unwrap_or(false);
4428 if !deferred {
4429 self.delete_pending_evict(&sub_id);
4430 }
4431 }
4432 self.rebuild_overlay_if_dirty();
4433 }
4434
4435 fn pinned_row_ids(&self, table: &str) -> std::collections::HashSet<String> {
4438 let mut pinned = std::collections::HashSet::new();
4439 for commit in &self.outbox {
4440 for op in &commit.ops {
4441 if op.table == table {
4442 pinned.insert(op.row_id.clone());
4443 }
4444 }
4445 }
4446 pinned
4447 }
4448
4449 fn save_pending_evict(&self, sub_id: &str, table: &str, effective: &[(String, Vec<String>)]) {
4450 let _ = self.conn.execute(
4451 "INSERT OR REPLACE INTO _syncular_window_pending_evict(sub_id, tbl, effective_scopes)
4452 VALUES (?1, ?2, ?3)",
4453 rusqlite::params![sub_id, table, scope_map_to_json(effective).to_string()],
4454 );
4455 }
4456
4457 fn delete_pending_evict(&self, sub_id: &str) {
4458 let _ = self.conn.execute(
4459 "DELETE FROM _syncular_window_pending_evict WHERE sub_id = ?1",
4460 rusqlite::params![sub_id],
4461 );
4462 }
4463
4464 fn load_pending_evictions(&self) -> Vec<PendingEvict> {
4465 let mut stmt = match self
4466 .conn
4467 .prepare("SELECT sub_id, tbl, effective_scopes FROM _syncular_window_pending_evict")
4468 {
4469 Ok(stmt) => stmt,
4470 Err(_) => return Vec::new(),
4471 };
4472 let rows = stmt.query_map([], |row| {
4473 Ok((
4474 row.get::<_, String>(0)?,
4475 row.get::<_, String>(1)?,
4476 row.get::<_, String>(2)?,
4477 ))
4478 });
4479 let mut out = Vec::new();
4480 if let Ok(rows) = rows {
4481 for entry in rows.filter_map(Result::ok) {
4482 let (sub_id, table, json) = entry;
4483 if let Ok(value) = serde_json::from_str::<Value>(&json) {
4484 if let Ok(effective) = json_to_scope_map(&value) {
4485 out.push((sub_id, table, effective));
4486 }
4487 }
4488 }
4489 }
4490 out
4491 }
4492
4493 pub fn mutate(&mut self, mutations: Vec<Mutation>) -> Result<String, String> {
4495 if mutations.is_empty() {
4496 return Err("a commit must contain at least one operation (§6.1)".to_owned());
4497 }
4498 let mut ops = Vec::with_capacity(mutations.len());
4499 for mutation in mutations {
4500 match mutation {
4501 Mutation::Upsert {
4502 table,
4503 values,
4504 base_version,
4505 } => {
4506 let schema_table = self
4507 .schema
4508 .table(&table)
4509 .ok_or_else(|| format!("unknown table {table:?}"))?;
4510 let values = normalize_values_casing(schema_table, values)?;
4514 let row_id = render_row_id_json(values.get(&schema_table.primary_key))?;
4515 encode_row_json(schema_table, &row_id, &values, &self.encryption)?;
4518 ops.push(OutboxOp {
4519 upsert: true,
4520 table,
4521 row_id,
4522 base_version,
4523 values: Some(values),
4524 changed_fields: None,
4525 });
4526 }
4527 Mutation::Delete {
4528 table,
4529 row_id,
4530 base_version,
4531 } => {
4532 if self.schema.table(&table).is_none() {
4533 return Err(format!("unknown table {table:?}"));
4534 }
4535 ops.push(OutboxOp {
4536 upsert: false,
4537 table,
4538 row_id,
4539 base_version,
4540 values: None,
4541 changed_fields: None,
4542 });
4543 }
4544 }
4545 }
4546 self.record_outbox_commit(ops)
4547 }
4548
4549 fn record_outbox_commit(&mut self, ops: Vec<OutboxOp>) -> Result<String, String> {
4550 let commit = OutboxCommit {
4551 client_commit_id: uuid::Uuid::new_v4().to_string(),
4552 ops,
4553 };
4554 self.begin_observation("syncular_mutation")?;
4555 let mut batch = ChangeAccumulator::default();
4556 for op in &commit.ops {
4557 let mut precise = self.record_row_scopes(&mut batch, &op.table, &op.row_id, false);
4558 if let Some(values) = &op.values {
4559 if let Some(table) = self.schema.table(&op.table) {
4560 for scope in &table.scope_variables {
4561 if let Some(Value::String(value)) = values.get(&scope.column) {
4562 batch.scope(&op.table, format!("{}:{value}", scope.prefix));
4563 precise = true;
4564 }
4565 }
4566 }
4567 }
4568 if !precise {
4569 batch.table(&op.table);
4570 }
4571 }
4572 self.persist_outbox_insert(&commit);
4573 let id = commit.client_commit_id.clone();
4574 self.outbox.push(commit);
4575 self.overlay_dirty.set(true);
4576 self.rebuild_overlay();
4577 batch.status = true;
4578 if let Err(error) = self.finish_observation("syncular_mutation", batch) {
4579 self.rollback_observation("syncular_mutation");
4580 return Err(error);
4581 }
4582 Ok(id)
4583 }
4584
4585 pub fn patch(
4590 &mut self,
4591 table: &str,
4592 row_id: &str,
4593 partial: Map<String, Value>,
4594 base_version: Option<i64>,
4595 ) -> Result<String, String> {
4596 let schema_table = self
4597 .schema
4598 .table(table)
4599 .ok_or_else(|| format!("unknown table {table:?}"))?;
4600 let partial = normalize_values_casing(schema_table, partial)?;
4601 let mut values = self
4602 .read_rows(table)?
4603 .into_iter()
4604 .find(|row| row.row_id == row_id)
4605 .map(|row| row.values)
4606 .ok_or_else(|| {
4607 format!(
4608 "sync.invalid_request: table {table:?} has no local row with primary key {row_id:?} to patch"
4609 )
4610 })?;
4611 let mut changed_fields = partial.keys().cloned().collect::<Vec<_>>();
4612 changed_fields.sort();
4613 values.extend(partial);
4614 let row_id_from_values = render_row_id_json(values.get(&schema_table.primary_key))?;
4615 if row_id_from_values != row_id {
4616 return Err("sync.invalid_request: patch cannot change the primary key".to_owned());
4617 }
4618 encode_row_json(schema_table, row_id, &values, &self.encryption)?;
4619 self.record_outbox_commit(vec![OutboxOp {
4620 upsert: true,
4621 table: table.to_owned(),
4622 row_id: row_id.to_owned(),
4623 base_version,
4624 values: Some(values),
4625 changed_fields: Some(changed_fields),
4626 }])
4627 }
4628
4629 pub fn pending_commit_ids(&self) -> Vec<String> {
4630 self.outbox
4631 .iter()
4632 .map(|c| c.client_commit_id.clone())
4633 .collect()
4634 }
4635
4636 pub fn conflicts(&self) -> &[ConflictRecord] {
4637 &self.conflicts
4638 }
4639
4640 pub fn rejections(&self) -> &[RejectionRecord] {
4641 &self.rejections
4642 }
4643
4644 pub fn resolve_commit_outcome(
4645 &mut self,
4646 input: ResolveCommitOutcomeInput,
4647 ) -> Result<CommitOutcome, String> {
4648 let current = self
4649 .commit_outcome(&input.client_commit_id)?
4650 .ok_or_else(|| {
4651 format!(
4652 "sync.outcome_not_found: no durable outcome exists for {:?}",
4653 input.client_commit_id
4654 )
4655 })?;
4656 if current.resolution != CommitOutcomeResolution::Active {
4657 return Ok(current);
4658 }
4659 if input.resolution == CommitOutcomeResolution::Active {
4660 return Err("sync.invalid_request: resolution must leave active state".to_owned());
4661 }
4662 match input.resolution {
4663 CommitOutcomeResolution::Superseded => {
4664 let replacement = input
4665 .replacement_client_commit_id
4666 .as_deref()
4667 .filter(|value| !value.is_empty() && *value != input.client_commit_id)
4668 .ok_or_else(|| {
4669 "sync.invalid_request: superseded outcomes require a distinct replacementClientCommitId"
4670 .to_owned()
4671 })?;
4672 let _ = replacement;
4673 }
4674 _ if input.replacement_client_commit_id.is_some() => {
4675 return Err(
4676 "sync.invalid_request: replacementClientCommitId is valid only for superseded outcomes"
4677 .to_owned(),
4678 );
4679 }
4680 _ => {}
4681 }
4682 let allowed = match current.status {
4683 CommitOutcomeStatus::Conflict => matches!(
4684 input.resolution,
4685 CommitOutcomeResolution::ResolvedKeepServer | CommitOutcomeResolution::Superseded
4686 ),
4687 CommitOutcomeStatus::Rejected => {
4688 input.resolution == CommitOutcomeResolution::Superseded
4689 }
4690 CommitOutcomeStatus::Applied | CommitOutcomeStatus::Cached => {
4691 input.resolution == CommitOutcomeResolution::Dismissed
4692 }
4693 };
4694 if !allowed {
4695 return Err(format!(
4696 "sync.invalid_request: resolution {:?} is invalid for {:?} outcome",
4697 input.resolution, current.status
4698 ));
4699 }
4700
4701 self.begin_observation("syncular_outcome_resolution")?;
4702 let result = (|| {
4703 self.conn
4704 .execute(
4705 "UPDATE _syncular_commit_outcomes
4706 SET resolution = ?1, resolved_at_ms = ?2,
4707 replacement_client_commit_id = ?3
4708 WHERE client_commit_id = ?4 AND resolution = 'active'",
4709 rusqlite::params![
4710 Self::outcome_resolution_name(input.resolution),
4711 self.clock_now_ms(),
4712 input.replacement_client_commit_id,
4713 input.client_commit_id
4714 ],
4715 )
4716 .map_err(|error| error.to_string())?;
4717 let resolved = self
4718 .commit_outcome(¤t.client_commit_id)?
4719 .ok_or_else(|| "sync.outcome_not_found: outcome disappeared".to_owned())?;
4720 self.prune_commit_outcomes()?;
4721 Ok(resolved)
4722 })();
4723 let resolved = match result {
4724 Ok(outcome) => outcome,
4725 Err(error) => {
4726 self.rollback_observation("syncular_outcome_resolution");
4727 return Err(error);
4728 }
4729 };
4730 self.conflicts
4731 .retain(|record| record.client_commit_id != current.client_commit_id);
4732 self.rejections
4733 .retain(|record| record.client_commit_id != current.client_commit_id);
4734 let batch = ChangeAccumulator {
4735 conflicts: current.status == CommitOutcomeStatus::Conflict,
4736 rejections: current.status == CommitOutcomeStatus::Rejected,
4737 outcomes: true,
4738 ..ChangeAccumulator::default()
4739 };
4740 if let Err(error) = self.finish_observation("syncular_outcome_resolution", batch) {
4741 self.rollback_observation("syncular_outcome_resolution");
4742 return Err(error);
4743 }
4744 Ok(resolved)
4745 }
4746
4747 pub fn schema_floor(&self) -> Option<&SchemaFloor> {
4748 self.schema_floor.as_ref()
4749 }
4750
4751 pub fn lease_state(&self) -> Option<&LeaseState> {
4753 self.lease_state.as_ref()
4754 }
4755
4756 fn record_lease_error(&mut self, code: &str) {
4759 if code != "sync.auth_lease_required" && code != "sync.auth_lease_revoked" {
4760 return;
4761 }
4762 let mut next = self.lease_state.clone().unwrap_or_default();
4763 next.error_code = Some(code.to_owned());
4764 self.set_lease_state(Some(next));
4765 }
4766
4767 fn set_lease_state(&mut self, next: Option<LeaseState>) {
4768 if self.lease_state == next {
4769 return;
4770 }
4771 if self.begin_observation("syncular_lease").is_err() {
4772 return;
4773 }
4774 self.lease_state = next;
4775 if let Some(lease) = &self.lease_state {
4776 if let Ok(json) = serde_json::to_string(lease) {
4777 self.set_meta(LEASE_STATE_KEY, &json);
4778 }
4779 } else {
4780 self.delete_meta(LEASE_STATE_KEY);
4781 }
4782 let batch = ChangeAccumulator {
4783 status: true,
4784 ..ChangeAccumulator::default()
4785 };
4786 if self.finish_observation("syncular_lease", batch).is_err() {
4787 self.rollback_observation("syncular_lease");
4788 }
4789 }
4790
4791 fn set_schema_floor(&mut self, next: Option<SchemaFloor>) {
4792 if self.schema_floor == next {
4793 return;
4794 }
4795 if self.begin_observation("syncular_schema_floor").is_err() {
4796 return;
4797 }
4798 self.schema_floor = next;
4799 self.stopped = self.schema_floor.is_some();
4800 if let Some(floor) = &self.schema_floor {
4801 if let Ok(json) = serde_json::to_string(floor) {
4802 self.set_meta(SCHEMA_FLOOR_KEY, &json);
4803 }
4804 } else {
4805 self.delete_meta(SCHEMA_FLOOR_KEY);
4806 }
4807 let batch = ChangeAccumulator {
4808 status: true,
4809 ..ChangeAccumulator::default()
4810 };
4811 if self
4812 .finish_observation("syncular_schema_floor", batch)
4813 .is_err()
4814 {
4815 self.rollback_observation("syncular_schema_floor");
4816 }
4817 }
4818
4819 fn set_upgrading(&mut self, value: bool) {
4820 if self.upgrading == value {
4821 return;
4822 }
4823 if self.begin_observation("syncular_upgrading").is_err() {
4824 return;
4825 }
4826 self.upgrading = value;
4827 let batch = ChangeAccumulator {
4828 status: true,
4829 ..ChangeAccumulator::default()
4830 };
4831 if self
4832 .finish_observation("syncular_upgrading", batch)
4833 .is_err()
4834 {
4835 self.rollback_observation("syncular_upgrading");
4836 }
4837 }
4838
4839 pub fn sync_needed(&self) -> bool {
4840 self.sync_needed
4841 }
4842
4843 pub fn subscription_state(&self, id: &str) -> Option<SubscriptionStateView> {
4844 let sub = self.subs.iter().find(|s| s.id == id)?;
4845 Some(SubscriptionStateView {
4846 id: sub.id.clone(),
4847 table: sub.table.clone(),
4848 status: sub.state.name().to_owned(),
4849 cursor: sub.cursor,
4850 has_resume_token: sub.bootstrap_state.is_some(),
4851 effective_scopes: sub.effective.as_ref().map(|e| scope_map_to_json(e)),
4852 reason_code: sub.reason_code.clone(),
4853 })
4854 }
4855
4856 pub fn read_rows(&self, table: &str) -> Result<Vec<RowState>, String> {
4857 let schema_table = self
4858 .schema
4859 .table(table)
4860 .ok_or_else(|| format!("unknown table {table:?}"))?;
4861 let sql = format!(
4862 "SELECT * FROM {} ORDER BY {} ASC",
4863 visible_table(table),
4864 quote_ident(&schema_table.primary_key)
4865 );
4866 let mut stmt = self.conn.prepare(&sql).map_err(|e| e.to_string())?;
4867 let mut rows = stmt.query([]).map_err(|e| e.to_string())?;
4868 let mut out = Vec::new();
4869 while let Some(row) = rows.next().map_err(|e| e.to_string())? {
4870 let mut values = Map::new();
4871 for (i, column) in schema_table.columns.iter().enumerate() {
4872 let value = row.get_ref(i).map_err(|e| e.to_string())?;
4873 values.insert(column.name.clone(), sql_ref_to_json(column, value));
4874 }
4875 let version: i64 = row
4876 .get(schema_table.columns.len())
4877 .map_err(|e| e.to_string())?;
4878 let row_id = match values.get(&schema_table.primary_key) {
4879 Some(Value::String(s)) => s.clone(),
4880 Some(Value::Number(n)) => n.to_string(),
4881 Some(Value::Bool(b)) => b.to_string(),
4882 other => format!("{}", other.cloned().unwrap_or(Value::Null)),
4883 };
4884 out.push(RowState {
4885 row_id,
4886 version,
4887 values,
4888 });
4889 }
4890 Ok(out)
4891 }
4892
4893 #[cfg(feature = "crdt-yjs")]
4909 fn crdt_column_bytes(
4910 &self,
4911 table: &str,
4912 row_id: &str,
4913 column: &str,
4914 ) -> Result<Option<Vec<u8>>, String> {
4915 let schema_table = self
4916 .schema
4917 .table(table)
4918 .ok_or_else(|| format!("unknown table {table:?}"))?;
4919 let col = schema_table
4920 .columns
4921 .iter()
4922 .find(|c| c.name == column)
4923 .ok_or_else(|| format!("table {table:?} has no column {column:?}"))?;
4924 if col.ty != ColumnType::Crdt {
4925 return Err(format!("column {column:?} is not a crdt column (§5.10.1)"));
4926 }
4927 let sql = format!(
4928 "SELECT {} FROM {} WHERE CAST({} AS TEXT) = ?1",
4929 quote_ident(column),
4930 visible_table(table),
4931 quote_ident(&schema_table.primary_key)
4932 );
4933 let bytes: Option<Vec<u8>> = self
4934 .conn
4935 .query_row(&sql, rusqlite::params![row_id], |row| {
4936 row.get::<_, Option<Vec<u8>>>(0)
4937 })
4938 .map_err(|e| match e {
4939 rusqlite::Error::QueryReturnedNoRows => "no such row".to_owned(),
4940 other => other.to_string(),
4941 })?;
4942 Ok(bytes)
4943 }
4944
4945 #[cfg(feature = "crdt-yjs")]
4949 pub fn crdt_text(
4950 &self,
4951 table: &str,
4952 row_id: &str,
4953 column: &str,
4954 name: &str,
4955 ) -> Result<String, String> {
4956 let bytes = self
4957 .crdt_column_bytes(table, row_id, column)?
4958 .unwrap_or_default();
4959 crate::crdt::text(&bytes, name)
4960 }
4961
4962 #[cfg(feature = "crdt-yjs")]
4966 pub fn crdt_insert_text(
4967 &mut self,
4968 table: &str,
4969 row_id: &str,
4970 column: &str,
4971 name: &str,
4972 index: u32,
4973 value: &str,
4974 ) -> Result<String, String> {
4975 let current = self
4976 .crdt_column_bytes(table, row_id, column)?
4977 .unwrap_or_default();
4978 let update = crate::crdt::insert_text(¤t, name, index, value)?;
4979 self.crdt_push_update(table, row_id, column, &update)
4980 }
4981
4982 #[cfg(feature = "crdt-yjs")]
4985 pub fn crdt_delete_text(
4986 &mut self,
4987 table: &str,
4988 row_id: &str,
4989 column: &str,
4990 name: &str,
4991 index: u32,
4992 len: u32,
4993 ) -> Result<String, String> {
4994 let current = self
4995 .crdt_column_bytes(table, row_id, column)?
4996 .unwrap_or_default();
4997 let update = crate::crdt::delete_text(¤t, name, index, len)?;
4998 self.crdt_push_update(table, row_id, column, &update)
4999 }
5000
5001 #[cfg(feature = "crdt-yjs")]
5006 pub fn crdt_apply_update(
5007 &mut self,
5008 table: &str,
5009 row_id: &str,
5010 column: &str,
5011 update: &[u8],
5012 ) -> Result<String, String> {
5013 let current = self
5014 .crdt_column_bytes(table, row_id, column)?
5015 .unwrap_or_default();
5016 let next = crate::crdt::apply_update(¤t, update)?;
5017 self.crdt_push_update(table, row_id, column, &next)
5018 }
5019
5020 #[cfg(feature = "crdt-yjs")]
5026 fn crdt_push_update(
5027 &mut self,
5028 table: &str,
5029 row_id: &str,
5030 column: &str,
5031 crdt_bytes: &[u8],
5032 ) -> Result<String, String> {
5033 let schema_table = self
5034 .schema
5035 .table(table)
5036 .ok_or_else(|| format!("unknown table {table:?}"))?
5037 .clone();
5038 let mut values: Map<String, Value> = self
5041 .read_rows(table)?
5042 .into_iter()
5043 .find(|r| r.row_id == row_id)
5044 .map(|r| r.values)
5045 .unwrap_or_else(|| {
5046 let mut map = Map::new();
5047 map.insert(
5048 schema_table.primary_key.clone(),
5049 Value::from(row_id.to_owned()),
5050 );
5051 map
5052 });
5053 let mut bytes_obj = Map::new();
5055 bytes_obj.insert("$bytes".to_owned(), Value::from(bytes_to_hex(crdt_bytes)));
5056 values.insert(column.to_owned(), Value::Object(bytes_obj));
5057 self.mutate(vec![Mutation::Upsert {
5058 table: table.to_owned(),
5059 values,
5060 base_version: None,
5061 }])
5062 }
5063
5064 pub fn query(&self, sql: &str, params: &[QueryValue]) -> Result<Vec<QueryRow>, String> {
5079 query_connection(&self.conn, sql, params)
5080 }
5081
5082 pub fn query_snapshot(
5084 &mut self,
5085 sql: &str,
5086 params: &[QueryValue],
5087 coverage: &[WindowCoverage],
5088 ) -> Result<QuerySnapshot, String> {
5089 snapshot_connection(&self.conn, sql, params, coverage)
5090 }
5091
5092 fn build_request(&self, url_capable: bool) -> (Message, RequestMeta) {
5095 let mut frames = vec![Frame::ReqHeader {
5096 client_id: self.client_id.clone(),
5097 schema_version: self.schema.version,
5098 }];
5099 let mut pushed_ids = Vec::new();
5100 let mut ops_in_request = 0usize;
5101 let mut deferred_commits = 0usize;
5102 for (index, commit) in self.outbox.iter().enumerate() {
5103 if ops_in_request > 0 && ops_in_request + commit.ops.len() > PUSH_OPS_PER_REQUEST {
5108 deferred_commits = self.outbox.len() - index;
5109 break;
5110 }
5111 ops_in_request += commit.ops.len();
5112 let operations = commit
5113 .ops
5114 .iter()
5115 .map(|op| {
5116 let payload = op.values.as_ref().and_then(|values| {
5117 let table = self.schema.table(&op.table)?;
5118 encode_row_json(table, &op.row_id, values, &self.encryption).ok()
5122 });
5123 ssp2::model::Operation {
5124 table: op.table.clone(),
5125 row_id: op.row_id.clone(),
5126 op: if op.upsert { Op::Upsert } else { Op::Delete },
5127 base_version: op.base_version,
5128 payload,
5129 }
5130 })
5131 .collect();
5132 frames.push(Frame::PushCommit {
5133 client_commit_id: commit.client_commit_id.clone(),
5134 operations,
5135 });
5136 pushed_ids.push(commit.client_commit_id.clone());
5137 }
5138 let accept = self.limits.accept.unwrap_or(if url_capable {
5141 DEFAULT_ACCEPT | ACCEPT_SIGNED_URLS
5142 } else {
5143 DEFAULT_ACCEPT
5144 });
5145 frames.push(Frame::PullHeader {
5146 limit_commits: self.limits.limit_commits.unwrap_or(0),
5147 limit_snapshot_rows: self.limits.limit_snapshot_rows.unwrap_or(0),
5148 max_snapshot_pages: self.limits.max_snapshot_pages.unwrap_or(0),
5149 accept,
5150 });
5151 let mut fresh = Vec::new();
5152 for sub in &self.subs {
5153 if sub.state != SubState::Active {
5154 continue;
5155 }
5156 let mut scopes = sub.requested.clone();
5157 sort_scope_map(&mut scopes);
5158 frames.push(Frame::Subscription {
5159 id: sub.id.clone(),
5160 table: sub.table.clone(),
5161 scopes,
5162 params: sub.params.clone().map(RawJson),
5163 cursor: sub.cursor,
5164 bootstrap_state: sub.bootstrap_state.clone().map(RawJson),
5165 });
5166 fresh.push((
5167 sub.id.clone(),
5168 sub.cursor < 0 && sub.bootstrap_state.is_none(),
5169 ));
5170 }
5171 let message = Message {
5172 msg_kind: MsgKind::Request,
5173 frames,
5174 };
5175 (
5176 message,
5177 RequestMeta {
5178 pushed_ids,
5179 fresh,
5180 accept,
5181 deferred_commits,
5182 },
5183 )
5184 }
5185
5186 pub fn sync(&mut self, transport: &mut dyn Transport) -> SyncOutcome {
5189 let started_at_ms = self.clock_now_ms();
5190 let outcome = self.sync_inner(transport);
5191 let completed_at_ms = self.clock_now_ms();
5192 self.last_round = Some(match &outcome {
5193 SyncOutcome::Ok(report) => DiagnosticLastRound {
5194 status: "succeeded".to_owned(),
5195 started_at_ms,
5196 completed_at_ms,
5197 duration_ms: completed_at_ms.saturating_sub(started_at_ms).max(0),
5198 counters: Some(DiagnosticRoundCounters {
5199 pushed: report.pushed,
5200 applied: report.applied.len(),
5201 rejected: report.rejected.len(),
5202 retryable: report.retryable.len(),
5203 conflicts: report.conflicts,
5204 commits_applied: report.commits_applied,
5205 segment_rows_applied: report.segment_rows_applied,
5206 bootstrapping: report.bootstrapping.len(),
5207 resets: report.resets.len(),
5208 revoked: report.revoked.len(),
5209 failed: report.failed.len(),
5210 deferred_commits: report.deferred_commits,
5211 }),
5212 error_code: None,
5213 },
5214 SyncOutcome::Failed { error_code, .. } => DiagnosticLastRound {
5215 status: "failed".to_owned(),
5216 started_at_ms,
5217 completed_at_ms,
5218 duration_ms: completed_at_ms.saturating_sub(started_at_ms).max(0),
5219 counters: None,
5220 error_code: Some(Self::diagnostic_code(error_code)),
5221 },
5222 });
5223 outcome
5224 }
5225
5226 fn sync_inner(&mut self, transport: &mut dyn Transport) -> SyncOutcome {
5227 if self.stopped {
5228 return SyncOutcome::Ok(SyncReport {
5231 schema_floor: self.schema_floor.clone(),
5232 ..SyncReport::default()
5233 });
5234 }
5235 self.set_sync_needed(false, false);
5238 if self.schema_has_blobs() {
5241 if let Err(TransportError { code, message }) = self.flush_blob_uploads(transport) {
5242 if Self::retryable_transport_code(&code) {
5243 self.schedule_background_retry();
5244 }
5245 return SyncOutcome::Failed {
5246 error_code: code,
5247 message,
5248 };
5249 }
5250 }
5251 let (message, meta) = self.build_request(transport.supports_url_fetch());
5252 let request_bytes = encode_message(&message);
5253 let round = if self.realtime_connected {
5257 transport.realtime_sync(&request_bytes)
5258 } else {
5259 transport.sync(&request_bytes)
5260 };
5261 let response_bytes = match round {
5262 Ok(bytes) => bytes,
5263 Err(TransportError { code, message }) => {
5264 self.record_lease_error(&code);
5267 if Self::retryable_transport_code(&code) {
5268 self.schedule_background_retry();
5269 }
5270 return SyncOutcome::Failed {
5271 error_code: code,
5272 message,
5273 };
5274 }
5275 };
5276 let response = match decode_message(&response_bytes) {
5277 Ok(message) => message,
5278 Err(error) => {
5279 return SyncOutcome::Failed {
5282 error_code: error.code.as_str().to_owned(),
5283 message: error.detail,
5284 };
5285 }
5286 };
5287 if response.msg_kind != MsgKind::Response {
5288 return SyncOutcome::Failed {
5289 error_code: "sync.invalid_request".to_owned(),
5290 message: "expected a response message".to_owned(),
5291 };
5292 }
5293 let mut outcome = self.process_response(transport, response, &meta);
5294 if let SyncOutcome::Ok(report) = &mut outcome {
5295 report.deferred_commits = meta.deferred_commits;
5296 }
5297 match &outcome {
5298 SyncOutcome::Ok(_) => self.reset_background_retry(),
5299 SyncOutcome::Failed { error_code, .. }
5300 if Self::retryable_transport_code(error_code) =>
5301 {
5302 self.schedule_background_retry();
5303 }
5304 SyncOutcome::Failed { .. } => {}
5305 }
5306 if meta.deferred_commits > 0 {
5307 self.set_sync_needed(true, true);
5310 }
5311 outcome
5312 }
5313
5314 pub fn sync_until_idle(
5315 &mut self,
5316 transport: &mut dyn Transport,
5317 max_rounds: Option<u32>,
5318 ) -> SyncOutcome {
5319 let rounds = max_rounds.unwrap_or(12).max(1);
5320 let mut aggregate = SyncReport::default();
5321 for _ in 0..rounds {
5322 match self.sync(transport) {
5323 SyncOutcome::Failed {
5324 error_code,
5325 message,
5326 } => {
5327 return SyncOutcome::Failed {
5328 error_code,
5329 message,
5330 };
5331 }
5332 SyncOutcome::Ok(report) => {
5333 aggregate.pushed += report.pushed;
5334 aggregate.applied.extend(report.applied.iter().cloned());
5335 aggregate.rejected.extend(report.rejected.iter().cloned());
5336 aggregate.retryable.extend(report.retryable.iter().cloned());
5337 aggregate.conflicts += report.conflicts;
5338 aggregate.commits_applied += report.commits_applied;
5339 aggregate.segment_rows_applied += report.segment_rows_applied;
5340 aggregate.bootstrapping = report.bootstrapping.clone();
5341 aggregate.resets.extend(report.resets.iter().cloned());
5342 aggregate.revoked.extend(report.revoked.iter().cloned());
5343 aggregate.failed.extend(report.failed.iter().cloned());
5344 aggregate.deferred_commits = report.deferred_commits;
5345 if report.schema_floor.is_some() {
5346 aggregate.schema_floor = report.schema_floor.clone();
5347 }
5348 let more = !report.bootstrapping.is_empty()
5354 || report.commits_applied > 0
5355 || report.segment_rows_applied > 0
5356 || !report.resets.is_empty()
5357 || self.sync_needed;
5358 if !more {
5359 break;
5360 }
5361 }
5362 }
5363 }
5364 SyncOutcome::Ok(aggregate)
5365 }
5366
5367 fn process_response(
5368 &mut self,
5369 transport: &mut dyn Transport,
5370 response: Message,
5371 meta: &RequestMeta,
5372 ) -> SyncOutcome {
5373 let mut report = SyncReport {
5374 pushed: meta.pushed_ids.len() as u32,
5375 ..SyncReport::default()
5376 };
5377 let mut rejection_details_by_commit: HashMap<String, BTreeMap<i32, RejectionDetails>> =
5378 HashMap::new();
5379 let pushed_ids = meta
5380 .pushed_ids
5381 .iter()
5382 .map(String::as_str)
5383 .collect::<HashSet<_>>();
5384 let mut last_final_push_result_id: Option<String> = None;
5385 for frame in &response.frames {
5386 match frame {
5387 Frame::PushResultDetails {
5388 client_commit_id,
5389 entries,
5390 } => {
5391 let details = rejection_details_by_commit
5392 .entry(client_commit_id.clone())
5393 .or_default();
5394 for entry in entries {
5395 let parsed = match RejectionDetails::parse(&entry.details.0) {
5396 Ok(value) => value,
5397 Err(message) => {
5398 return SyncOutcome::Failed {
5399 error_code: "sync.invalid_request".to_owned(),
5400 message,
5401 };
5402 }
5403 };
5404 details.insert(entry.op_index, parsed);
5405 }
5406 }
5407 Frame::PushResult {
5408 client_commit_id,
5409 status,
5410 results,
5411 ..
5412 } if pushed_ids.contains(client_commit_id.as_str())
5413 && Self::push_result_is_final(*status, results) =>
5414 {
5415 last_final_push_result_id = Some(client_commit_id.clone());
5416 }
5417 _ => {}
5418 }
5419 }
5420 let mut frames = response.frames.into_iter();
5421 match frames.next() {
5422 Some(Frame::RespHeader {
5423 required_schema_version,
5424 latest_schema_version,
5425 }) => {
5426 if let Some(required) = required_schema_version {
5427 let floor = SchemaFloor {
5430 required_schema_version: Some(required),
5431 latest_schema_version,
5432 };
5433 self.set_schema_floor(Some(floor.clone()));
5434 report.schema_floor = Some(floor);
5435 return SyncOutcome::Ok(report);
5436 }
5437 }
5438 _ => {
5439 return SyncOutcome::Failed {
5440 error_code: "sync.invalid_request".to_owned(),
5441 message: "response does not start with RESP_HEADER".to_owned(),
5442 };
5443 }
5444 }
5445
5446 let mut failure: Option<(String, String)> = None;
5447 while let Some(frame) = frames.next() {
5448 match frame {
5449 Frame::PushResult {
5450 client_commit_id,
5451 status,
5452 commit_seq: _,
5453 results,
5454 } => {
5455 let prune_outcomes =
5456 last_final_push_result_id.as_deref() == Some(&client_commit_id);
5457 self.handle_push_result(
5458 &client_commit_id,
5459 status,
5460 &results,
5461 rejection_details_by_commit.get(&client_commit_id),
5462 &mut report,
5463 prune_outcomes,
5464 );
5465 }
5466 Frame::PushResultDetails { .. } => {}
5467 Frame::SubStart {
5468 id,
5469 status,
5470 reason_code,
5471 effective_scopes,
5472 bootstrap: _,
5473 } => {
5474 let mut body = Vec::new();
5475 let mut sub_end: Option<(i64, Option<String>)> = None;
5476 for inner in frames.by_ref() {
5477 match inner {
5478 Frame::SubEnd {
5479 next_cursor,
5480 bootstrap_state,
5481 } => {
5482 sub_end = Some((next_cursor, bootstrap_state.map(|r| r.0)));
5483 break;
5484 }
5485 Frame::Unknown { .. } => {}
5486 other => body.push(other),
5487 }
5488 }
5489 let Some((next_cursor, bootstrap_state)) = sub_end else {
5490 failure = Some((
5491 "sync.invalid_request".to_owned(),
5492 "subscription section without SUB_END".to_owned(),
5493 ));
5494 break;
5495 };
5496 if let Err(SectionError::Abort(code, message)) = self.process_section(
5497 transport,
5498 &id,
5499 status,
5500 &reason_code,
5501 effective_scopes,
5502 body,
5503 next_cursor,
5504 bootstrap_state,
5505 meta,
5506 &mut report,
5507 ) {
5508 failure = Some((code, message));
5509 break;
5510 }
5511 }
5512 Frame::Lease {
5513 lease_id,
5514 expires_at_ms,
5515 } => {
5516 self.set_lease_state(Some(LeaseState {
5519 lease_id: Some(lease_id),
5520 expires_at_ms: Some(expires_at_ms),
5521 error_code: None,
5522 }));
5523 }
5524 Frame::Error { code, message, .. } => {
5525 failure = Some((code, message));
5528 break;
5529 }
5530 Frame::Unknown { .. } => {}
5531 _ => {
5532 failure = Some((
5533 "sync.invalid_request".to_owned(),
5534 "unexpected frame in response".to_owned(),
5535 ));
5536 break;
5537 }
5538 }
5539 }
5540
5541 if self.overlay_dirty.get() {
5546 self.rebuild_overlay();
5547 }
5548 self.reconcile_blob_refcounts(false);
5555
5556 if let Some((error_code, message)) = failure {
5557 return SyncOutcome::Failed {
5558 error_code,
5559 message,
5560 };
5561 }
5562 self.drain_pending_evictions();
5565 self.ack_after_pull(transport);
5566 if self.upgrading && report.bootstrapping.is_empty() {
5569 self.set_upgrading(false);
5570 }
5571 SyncOutcome::Ok(report)
5572 }
5573
5574 fn handle_push_result(
5577 &mut self,
5578 client_commit_id: &str,
5579 status: PushStatus,
5580 results: &[OpResult],
5581 rejection_details: Option<&BTreeMap<i32, RejectionDetails>>,
5582 report: &mut SyncReport,
5583 prune_outcomes: bool,
5584 ) {
5585 let Some(index) = self
5586 .outbox
5587 .iter()
5588 .position(|c| c.client_commit_id == client_commit_id)
5589 else {
5590 return;
5591 };
5592 if self.begin_observation("syncular_push_result").is_err() {
5593 return;
5594 }
5595 let mut batch = ChangeAccumulator::default();
5596 let operations = self.outbox[index].ops.clone();
5597 match status {
5598 PushStatus::Applied | PushStatus::Cached => {
5599 let journal_results = results
5602 .iter()
5603 .map(|result| {
5604 let op_index = match result {
5605 OpResult::Applied { op_index }
5606 | OpResult::Conflict { op_index, .. }
5607 | OpResult::Error { op_index, .. } => *op_index,
5608 };
5609 CommitOperationOutcome::Applied { op_index }
5610 })
5611 .collect::<Vec<_>>();
5612 let outcome_status = if status == PushStatus::Applied {
5613 CommitOutcomeStatus::Applied
5614 } else {
5615 CommitOutcomeStatus::Cached
5616 };
5617 let persisted = self
5618 .persist_commit_outcome(
5619 client_commit_id,
5620 outcome_status,
5621 &journal_results,
5622 None,
5623 )
5624 .and_then(|()| self.delete_outbox_persisted(client_commit_id));
5625 let persisted = if prune_outcomes {
5626 persisted.and_then(|()| self.prune_commit_outcomes())
5627 } else {
5628 persisted
5629 };
5630 if persisted.is_err() {
5631 self.rollback_observation("syncular_push_result");
5632 return;
5633 }
5634 report.applied.push(client_commit_id.to_owned());
5635 self.outbox.remove(index);
5636 self.overlay_dirty.set(true);
5637 batch.status = true;
5638 batch.outcomes = true;
5639 }
5640 PushStatus::Rejected => {
5641 if results.iter().any(|result| {
5642 matches!(
5643 result,
5644 OpResult::Error {
5645 code,
5646 retryable: true,
5647 ..
5648 } if code == "sync.idempotency_cache_miss"
5649 )
5650 }) {
5651 report.retryable.push(client_commit_id.to_owned());
5654 if self
5655 .finish_observation("syncular_push_result", batch)
5656 .is_err()
5657 {
5658 self.rollback_observation("syncular_push_result");
5659 }
5660 return;
5661 }
5662
5663 let mut journal_results = Vec::with_capacity(results.len());
5664 let mut conflicts = Vec::new();
5665 let mut rejections = Vec::new();
5666 for result in results {
5667 match result {
5668 OpResult::Applied { op_index } => {
5669 journal_results.push(CommitOperationOutcome::Applied {
5670 op_index: *op_index,
5671 });
5672 }
5673 OpResult::Conflict {
5674 op_index,
5675 code,
5676 message,
5677 server_version,
5678 server_row,
5679 } => {
5680 let operation = operations
5681 .get(*op_index as usize)
5682 .map(CommitOperation::from);
5683 let (table, row_id) = operation
5684 .as_ref()
5685 .map(|op| (op.table.clone(), op.row_id.clone()))
5686 .unwrap_or_default();
5687 let server_row_json = self
5688 .schema
5689 .table(&table)
5690 .and_then(|t| {
5691 decode_row_bytes(t, server_row, &self.encryption)
5692 .ok()
5693 .map(|row| (t, row))
5694 })
5695 .map(|(t, row)| {
5696 let mut map = Map::new();
5697 for (i, column) in t.columns.iter().enumerate() {
5698 map.insert(
5699 column.name.clone(),
5700 column_value_to_json(row.get(i).unwrap_or(&None)),
5701 );
5702 }
5703 map
5704 })
5705 .unwrap_or_default();
5706 let conflict = ConflictRecord {
5707 client_commit_id: client_commit_id.to_owned(),
5708 op_index: *op_index,
5709 table,
5710 row_id,
5711 code: code.clone(),
5712 message: message.clone(),
5713 server_version: *server_version,
5714 server_row: server_row_json,
5715 operation,
5716 };
5717 journal_results.push(CommitOperationOutcome::Conflict {
5718 conflict: conflict.clone(),
5719 });
5720 conflicts.push(conflict);
5721 }
5722 OpResult::Error {
5723 op_index,
5724 code,
5725 message,
5726 retryable,
5727 } => {
5728 let rejection = RejectionRecord {
5729 client_commit_id: client_commit_id.to_owned(),
5730 op_index: *op_index,
5731 code: code.clone(),
5732 message: message.clone(),
5733 retryable: *retryable,
5734 details: rejection_details
5735 .and_then(|details| details.get(op_index))
5736 .cloned(),
5737 operation: operations
5738 .get(*op_index as usize)
5739 .map(CommitOperation::from),
5740 };
5741 journal_results.push(CommitOperationOutcome::Error {
5742 rejection: rejection.clone(),
5743 });
5744 rejections.push(rejection);
5745 }
5746 }
5747 }
5748 let outcome_status = if conflicts.is_empty() {
5749 CommitOutcomeStatus::Rejected
5750 } else {
5751 CommitOutcomeStatus::Conflict
5752 };
5753 let persisted = self
5754 .persist_commit_outcome(
5755 client_commit_id,
5756 outcome_status,
5757 &journal_results,
5758 Some(&operations),
5759 )
5760 .and_then(|()| self.delete_outbox_persisted(client_commit_id));
5761 let persisted = if prune_outcomes {
5762 persisted.and_then(|()| self.prune_commit_outcomes())
5763 } else {
5764 persisted
5765 };
5766 if persisted.is_err() {
5767 self.rollback_observation("syncular_push_result");
5768 return;
5769 }
5770 report.conflicts += conflicts.len() as u32;
5771 report.rejected.push(client_commit_id.to_owned());
5772 batch.conflicts = !conflicts.is_empty();
5773 batch.rejections = !rejections.is_empty();
5774 batch.status = true;
5775 batch.outcomes = true;
5776 self.conflicts.extend(conflicts);
5777 self.rejections.extend(rejections);
5778 self.outbox.remove(index);
5779 self.overlay_dirty.set(true);
5780 }
5781 }
5782 if batch.status {
5783 for operation in &operations {
5784 if !self.record_row_scopes(&mut batch, &operation.table, &operation.row_id, false) {
5785 batch.table(&operation.table);
5786 }
5787 }
5788 }
5789 if self
5790 .finish_observation("syncular_push_result", batch)
5791 .is_err()
5792 {
5793 self.rollback_observation("syncular_push_result");
5794 }
5795 }
5796
5797 fn push_result_is_final(status: PushStatus, results: &[OpResult]) -> bool {
5798 status != PushStatus::Rejected
5799 || !results.iter().any(|result| {
5800 matches!(
5801 result,
5802 OpResult::Error {
5803 code,
5804 retryable: true,
5805 ..
5806 } if code == "sync.idempotency_cache_miss"
5807 )
5808 })
5809 }
5810
5811 #[allow(clippy::too_many_arguments)]
5814 fn process_section(
5815 &mut self,
5816 transport: &mut dyn Transport,
5817 id: &str,
5818 status: SubStatus,
5819 reason_code: &str,
5820 effective_scopes: Vec<(String, Vec<String>)>,
5821 body: Vec<Frame>,
5822 next_cursor: i64,
5823 bootstrap_state: Option<String>,
5824 meta: &RequestMeta,
5825 report: &mut SyncReport,
5826 ) -> Result<(), SectionError> {
5827 let Some(sub_index) = self.subs.iter().position(|s| s.id == id) else {
5828 return Ok(()); };
5830 match status {
5831 SubStatus::Revoked => {
5832 self.begin_observation("syncular_revocation")
5833 .map_err(|message| SectionError::Abort("storage.failed".to_owned(), message))?;
5834 let mut batch = ChangeAccumulator::default();
5835 let registered = self.window_unit_by_sub_id(id);
5836 let (table, effective) = {
5838 let sub = &self.subs[sub_index];
5839 (sub.table.clone(), sub.effective.clone().unwrap_or_default())
5840 };
5841 let purged = self.purge_scope_rows(&table, &effective);
5842 match purged {
5843 Ok(()) => {
5844 self.record_scope_map(&mut batch, &table, &effective);
5845 let sub = &mut self.subs[sub_index];
5846 sub.state = SubState::Revoked;
5847 sub.reason_code = Some(if reason_code.is_empty() {
5848 "sync.scope_revoked".to_owned()
5849 } else {
5850 reason_code.to_owned()
5851 });
5852 report.revoked.push(id.to_owned());
5853 let doomed_effective = effective;
5854 let sub_table = table;
5855 self.persist_sub(&self.subs[sub_index].clone());
5856 let dropped = self
5857 .drop_doomed_outbox(&sub_table, &doomed_effective)
5858 .map_err(|message| {
5859 SectionError::Abort("storage.failed".to_owned(), message)
5860 })?;
5861 if dropped {
5862 batch.status = true;
5863 batch.rejections = true;
5864 batch.outcomes = true;
5865 }
5866 self.reconcile_blob_refcounts(true);
5869 }
5870 Err(()) => {
5871 let sub = &mut self.subs[sub_index];
5874 sub.state = SubState::Failed;
5875 sub.reason_code = Some("sync.scope_revoked".to_owned());
5876 report.failed.push(id.to_owned());
5877 self.persist_sub(&self.subs[sub_index].clone());
5878 }
5879 }
5880 if let Some((base_key, unit)) = registered {
5881 batch.window(&base_key, &self.subs[sub_index].table, &unit);
5882 }
5883 self.rebuild_overlay_if_dirty();
5884 self.finish_observation("syncular_revocation", batch)
5885 .map_err(|message| SectionError::Abort("storage.failed".to_owned(), message))?;
5886 Ok(())
5887 }
5888 SubStatus::Reset => {
5889 self.begin_observation("syncular_reset")
5890 .map_err(|message| SectionError::Abort("storage.failed".to_owned(), message))?;
5891 let mut batch = ChangeAccumulator::default();
5892 let registered = self.window_unit_by_sub_id(id);
5893 let sub = &mut self.subs[sub_index];
5896 sub.cursor = -1;
5897 sub.bootstrap_state = None;
5898 report.resets.push(id.to_owned());
5899 self.persist_sub(&self.subs[sub_index].clone());
5900 if let Some((base_key, unit)) = registered {
5901 batch.window(&base_key, &self.subs[sub_index].table, &unit);
5902 }
5903 self.finish_observation("syncular_reset", batch)
5904 .map_err(|message| SectionError::Abort("storage.failed".to_owned(), message))?;
5905 Ok(())
5906 }
5907 SubStatus::Active => {
5908 let fresh = meta
5909 .fresh
5910 .iter()
5911 .find(|(fid, _)| fid == id)
5912 .map(|(_, f)| *f)
5913 .unwrap_or(false);
5914 let was_pending = self.subs[sub_index].cursor < 0
5915 || self.subs[sub_index].bootstrap_state.is_some();
5916 let registered = self.window_unit_by_sub_id(id);
5917 self.subs[sub_index].effective = Some(effective_scopes);
5919 self.begin_observation("syncular_section")
5920 .map_err(|message| SectionError::Abort("storage.failed".to_owned(), message))?;
5921 let mut batch = ChangeAccumulator::default();
5922 let outcome = self.apply_section_body(
5923 transport, sub_index, body, fresh, meta, report, &mut batch,
5924 );
5925 match outcome {
5926 Ok(()) => {
5927 let sub = &mut self.subs[sub_index];
5928 sub.cursor = next_cursor;
5931 sub.bootstrap_state = bootstrap_state;
5932 sub.synced_once = true;
5933 if sub.bootstrap_state.is_some() {
5934 report.bootstrapping.push(id.to_owned());
5935 }
5936 let completed =
5937 was_pending && sub.cursor >= 0 && sub.bootstrap_state.is_none();
5938 self.persist_sub(&self.subs[sub_index].clone());
5939 if completed {
5940 if let Some((base_key, unit)) = registered.clone() {
5941 batch.window(&base_key, &self.subs[sub_index].table, &unit);
5942 }
5943 }
5944 self.rebuild_overlay_if_dirty();
5945 self.finish_observation("syncular_section", batch)
5946 .map_err(|message| {
5947 SectionError::Abort("storage.failed".to_owned(), message)
5948 })?;
5949 Ok(())
5950 }
5951 Err(SectionError::FailClosed) => {
5952 self.rollback_observation("syncular_section");
5955 self.begin_observation("syncular_section_failure")
5956 .map_err(|message| {
5957 SectionError::Abort("storage.failed".to_owned(), message)
5958 })?;
5959 let mut failure_batch = ChangeAccumulator::default();
5960 let sub = &mut self.subs[sub_index];
5961 sub.state = SubState::Failed;
5962 sub.reason_code = Some("sync.scope_revoked".to_owned());
5963 report.failed.push(id.to_owned());
5964 self.persist_sub(&self.subs[sub_index].clone());
5965 if let Some((base_key, unit)) = registered {
5966 failure_batch.window(&base_key, &self.subs[sub_index].table, &unit);
5967 }
5968 self.finish_observation("syncular_section_failure", failure_batch)
5969 .map_err(|message| {
5970 SectionError::Abort("storage.failed".to_owned(), message)
5971 })?;
5972 Ok(())
5973 }
5974 Err(SectionError::Abort(code, message)) => {
5975 self.rollback_observation("syncular_section");
5978 Err(SectionError::Abort(code, message))
5979 }
5980 }
5981 }
5982 }
5983 }
5984
5985 #[allow(clippy::too_many_arguments)]
5989 fn apply_section_body(
5990 &mut self,
5991 transport: &mut dyn Transport,
5992 sub_index: usize,
5993 body: Vec<Frame>,
5994 fresh: bool,
5995 meta: &RequestMeta,
5996 report: &mut SyncReport,
5997 batch: &mut ChangeAccumulator,
5998 ) -> Result<(), SectionError> {
5999 let mut saw_segment = false;
6000 for frame in body {
6001 match frame {
6002 Frame::Commit {
6003 tables, changes, ..
6004 } => {
6005 self.record_commit_changes(batch, &tables, &changes);
6006 self.apply_commit_changes(&tables, &changes)
6007 .map_err(|(c, m)| SectionError::Abort(c, m))?;
6008 report.commits_applied += 1;
6009 }
6010 Frame::SegmentInline { payload } => {
6011 let segment = decode_rows_segment(&payload)
6012 .map_err(|e| SectionError::Abort(e.code.as_str().to_owned(), e.detail))?;
6013 let first = !saw_segment;
6014 saw_segment = true;
6015 let effective = self.subs[sub_index].effective.clone().unwrap_or_default();
6016 let cleared =
6017 fresh && first && self.scoped_rows_exist(&segment.table, &effective);
6018 let applied = self.apply_segment(sub_index, &segment, fresh && first)?;
6019 if applied > 0 || cleared {
6020 batch.table(&segment.table);
6021 }
6022 report.segment_rows_applied += applied;
6023 }
6024 Frame::SegmentRef {
6025 segment_id,
6026 media_type,
6027 table,
6028 row_count,
6029 as_of_commit_seq,
6030 scope_digest,
6031 row_cursor,
6032 next_row_cursor,
6033 url,
6034 url_expires_at_ms,
6035 ..
6036 } => {
6037 let advertised = match media_type {
6040 MediaType::Rows => {
6041 meta.accept & ACCEPT_EXTERNAL_ROWS != 0
6042 || meta.accept & ACCEPT_INLINE_ROWS != 0
6043 }
6044 MediaType::Sqlite => meta.accept & ACCEPT_SQLITE != 0,
6045 };
6046 if !advertised {
6047 return Err(SectionError::Abort(
6048 "sync.invalid_request".to_owned(),
6049 format!(
6050 "SEGMENT_REF mediaType {} was not advertised in accept (§4.2)",
6051 media_type.name()
6052 ),
6053 ));
6054 }
6055 let bytes = if let Some(url) = url {
6056 if meta.accept & ACCEPT_SIGNED_URLS == 0 {
6061 return Err(SectionError::Abort(
6062 "sync.invalid_request".to_owned(),
6063 "SEGMENT_REF carries a url but accept bit 3 was not advertised (§5.4)"
6064 .to_owned(),
6065 ));
6066 }
6067 if url_expires_at_ms.is_some_and(|exp| exp <= self.clock_now_ms()) {
6069 return Err(SectionError::Abort(
6070 "sync.segment_expired".to_owned(),
6071 format!(
6072 "signed URL for segment {segment_id} expired before fetch — re-pull mints fresh descriptors (§5.4)"
6073 ),
6074 ));
6075 }
6076 transport
6077 .fetch_url(&url)
6078 .map_err(|e| SectionError::Abort(e.code, e.message))?
6079 } else {
6080 let requested_scopes_json =
6081 canonical_scope_json(&self.subs[sub_index].requested);
6082 transport
6083 .download_segment(&SegmentRequest {
6084 segment_id: segment_id.clone(),
6085 table,
6086 requested_scopes_json,
6087 })
6088 .map_err(|e| SectionError::Abort(e.code, e.message))?
6089 };
6090 let digest = Sha256::digest(&bytes);
6092 let expected = segment_id
6093 .strip_prefix("sha256:")
6094 .unwrap_or(segment_id.as_str());
6095 if bytes_to_hex(&digest) != expected {
6096 return Err(SectionError::Abort(
6097 "sync.invalid_request".to_owned(),
6098 "segment bytes do not match the content address (§5.1)".to_owned(),
6099 ));
6100 }
6101 if media_type == MediaType::Sqlite {
6102 if row_cursor.is_some() || next_row_cursor.is_some() {
6105 return Err(SectionError::Abort(
6106 "sync.invalid_request".to_owned(),
6107 "sqlite segments are whole-table: rowCursor/nextRowCursor must be absent (§5.3)"
6108 .to_owned(),
6109 ));
6110 }
6111 let first = !saw_segment;
6112 saw_segment = true;
6113 let sub_table = self.subs[sub_index].table.clone();
6114 let effective = self.subs[sub_index].effective.clone().unwrap_or_default();
6115 let cleared =
6116 fresh && first && self.scoped_rows_exist(&sub_table, &effective);
6117 let applied = self.apply_sqlite_segment(
6118 sub_index,
6119 &bytes,
6120 fresh && first,
6121 row_count,
6122 as_of_commit_seq,
6123 &scope_digest,
6124 )?;
6125 if applied > 0 || cleared {
6126 batch.table(&sub_table);
6127 }
6128 report.segment_rows_applied += applied;
6129 } else {
6130 let segment = decode_rows_segment(&bytes).map_err(|e| {
6131 SectionError::Abort(e.code.as_str().to_owned(), e.detail)
6132 })?;
6133 let first = row_cursor.is_none();
6134 saw_segment = true;
6135 let effective = self.subs[sub_index].effective.clone().unwrap_or_default();
6136 let cleared =
6137 fresh && first && self.scoped_rows_exist(&segment.table, &effective);
6138 let applied = self.apply_segment(sub_index, &segment, fresh && first)?;
6139 if applied > 0 || cleared {
6140 batch.table(&segment.table);
6141 }
6142 report.segment_rows_applied += applied;
6143 }
6144 }
6145 Frame::Unknown { .. } => {}
6146 _ => {
6147 return Err(SectionError::Abort(
6148 "sync.invalid_request".to_owned(),
6149 "unexpected frame inside a subscription section".to_owned(),
6150 ));
6151 }
6152 }
6153 }
6154 Ok(())
6155 }
6156
6157 fn apply_commit_changes(
6158 &mut self,
6159 tables: &[String],
6160 changes: &[ssp2::model::Change],
6161 ) -> Result<(), (String, String)> {
6162 for change in changes {
6163 let table_name = tables.get(change.table_index as usize).ok_or_else(|| {
6164 (
6165 "sync.invalid_request".to_owned(),
6166 "change tableIndex out of range".to_owned(),
6167 )
6168 })?;
6169 let table = self.schema.table(table_name).ok_or_else(|| {
6170 (
6171 "sync.schema_mismatch".to_owned(),
6172 format!("change targets unknown table {table_name:?}"),
6173 )
6174 })?;
6175 match change.op {
6176 Op::Upsert => {
6177 let payload = change.row.as_ref().ok_or_else(|| {
6178 (
6179 "sync.invalid_request".to_owned(),
6180 "upsert change without row payload".to_owned(),
6181 )
6182 })?;
6183 let row = decode_row_bytes(table, payload, &self.encryption)
6185 .map_err(|m| ("sync.invalid_request".to_owned(), m))?;
6186 let version = change.row_version.unwrap_or(0);
6187 let table_name = table.name.clone();
6188 self.write_base_row(&table_name, &row, version)
6189 .map_err(|m| ("sync.invalid_request".to_owned(), m))?;
6190 }
6191 Op::Delete => {
6192 self.delete_base_row(table_name, &change.row_id)
6193 .map_err(|m| ("sync.invalid_request".to_owned(), m))?;
6194 }
6195 }
6196 }
6197 Ok(())
6198 }
6199
6200 fn apply_segment(
6205 &mut self,
6206 sub_index: usize,
6207 segment: &RowsSegment,
6208 first_fresh_page: bool,
6209 ) -> Result<u32, SectionError> {
6210 let (sub_table, effective) = {
6211 let sub = &self.subs[sub_index];
6212 (sub.table.clone(), sub.effective.clone().unwrap_or_default())
6213 };
6214 let table = self.schema.table(&sub_table).cloned().ok_or_else(|| {
6215 SectionError::Abort(
6216 "sync.schema_mismatch".to_owned(),
6217 format!("subscription table {sub_table:?} is not in the client schema"),
6218 )
6219 })?;
6220 let matches = segment.table == table.name
6225 && segment.schema_version == self.schema.version
6226 && segment.columns.len() == table.wire_columns.len()
6227 && segment
6228 .columns
6229 .iter()
6230 .zip(table.wire_columns.iter())
6231 .all(|(a, b)| a.name == b.name && a.ty == b.ty && a.nullable == b.nullable);
6232 if !matches {
6233 return Err(SectionError::Abort(
6234 "sync.schema_mismatch".to_owned(),
6235 "segment column table does not match the generated schema (§5.2)".to_owned(),
6236 ));
6237 }
6238 if first_fresh_page {
6239 self.purge_scope_rows(&table.name, &effective)
6243 .map_err(|()| SectionError::FailClosed)?;
6244 }
6245 let mut applied = 0u32;
6246 for block in &segment.blocks {
6247 for row in block {
6248 let decrypted;
6253 let values = if table.has_encrypted_columns() {
6254 let mut values = row.values.clone();
6255 crate::values::decrypt_segment_row(&table, &mut values, &self.encryption)
6256 .map_err(|m| SectionError::Abort("client.decrypt_failed".to_owned(), m))?;
6257 decrypted = values;
6258 &decrypted
6259 } else {
6260 &row.values
6261 };
6262 self.write_base_row(&table.name, values, row.server_version)
6265 .map_err(|m| SectionError::Abort("sync.invalid_request".to_owned(), m))?;
6266 applied += 1;
6267 }
6268 }
6269 Ok(applied)
6270 }
6271
6272 fn apply_sqlite_segment(
6280 &mut self,
6281 sub_index: usize,
6282 bytes: &[u8],
6283 first_fresh_page: bool,
6284 row_count: i64,
6285 as_of_commit_seq: i64,
6286 scope_digest: &str,
6287 ) -> Result<u32, SectionError> {
6288 let invalid = |detail: &str| {
6289 SectionError::Abort(
6290 "sync.invalid_request".to_owned(),
6291 format!("sqlite segment rejected: {detail} (§5.3)"),
6292 )
6293 };
6294 let (sub_table, effective) = {
6295 let sub = &self.subs[sub_index];
6296 (sub.table.clone(), sub.effective.clone().unwrap_or_default())
6297 };
6298 let table = self.schema.table(&sub_table).cloned().ok_or_else(|| {
6299 SectionError::Abort(
6300 "sync.schema_mismatch".to_owned(),
6301 format!("subscription table {sub_table:?} is not in the client schema"),
6302 )
6303 })?;
6304
6305 let path = std::env::temp_dir().join(format!("syncular-image-{}.db", uuid::Uuid::new_v4()));
6306 std::fs::write(&path, bytes).map_err(|_| invalid("image temp file write failed"))?;
6307 let img = match rusqlite::Connection::open_with_flags(
6308 &path,
6309 rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY,
6310 ) {
6311 Ok(conn) => conn,
6312 Err(_) => {
6313 let _ = std::fs::remove_file(&path);
6314 return Err(invalid("bytes do not open as a SQLite database"));
6315 }
6316 };
6317 let outcome = self.apply_sqlite_image(
6318 &img,
6319 &table,
6320 first_fresh_page,
6321 &effective,
6322 row_count,
6323 as_of_commit_seq,
6324 scope_digest,
6325 );
6326 drop(img);
6327 let _ = std::fs::remove_file(&path);
6328 outcome
6329 }
6330
6331 #[allow(clippy::too_many_arguments)]
6332 fn apply_sqlite_image(
6333 &mut self,
6334 img: &rusqlite::Connection,
6335 table: &crate::schema::TableSchema,
6336 first_fresh_page: bool,
6337 effective: &[(String, Vec<String>)],
6338 row_count: i64,
6339 as_of_commit_seq: i64,
6340 scope_digest: &str,
6341 ) -> Result<u32, SectionError> {
6342 let invalid = |detail: String| {
6343 SectionError::Abort(
6344 "sync.invalid_request".to_owned(),
6345 format!("sqlite segment rejected: {detail} (§5.3)"),
6346 )
6347 };
6348
6349 type MetaRow = (i64, String, i64, i64, String, i64, i64);
6352 let meta: MetaRow = img
6353 .query_row(
6354 "SELECT format, \"table\", \"schemaVersion\", \"asOfCommitSeq\",
6355 \"scopeDigest\", \"rowCount\",
6356 (SELECT count(*) FROM _syncular_segment)
6357 FROM _syncular_segment",
6358 [],
6359 |row| {
6360 Ok((
6361 row.get(0)?,
6362 row.get(1)?,
6363 row.get(2)?,
6364 row.get(3)?,
6365 row.get(4)?,
6366 row.get(5)?,
6367 row.get(6)?,
6368 ))
6369 },
6370 )
6371 .map_err(|_| invalid("missing or unreadable _syncular_segment metadata".to_owned()))?;
6372 let (format, meta_table, schema_version, pin, digest, meta_rows, meta_count) = meta;
6373 if meta_count != 1 {
6374 return Err(invalid(format!(
6375 "_syncular_segment must contain exactly one row, found {meta_count}"
6376 )));
6377 }
6378 if format != 1 {
6379 return Err(invalid(format!("format {format}")));
6380 }
6381 if meta_table != table.name {
6382 return Err(invalid(format!("image table {meta_table:?}")));
6383 }
6384 if schema_version != i64::from(self.schema.version) {
6385 return Err(invalid(format!("schemaVersion {schema_version}")));
6386 }
6387 if pin != as_of_commit_seq {
6388 return Err(invalid(format!("asOfCommitSeq {pin}")));
6389 }
6390 if digest != scope_digest {
6391 return Err(invalid("scopeDigest mismatch".to_owned()));
6392 }
6393 if meta_rows != row_count {
6394 return Err(invalid(format!("rowCount {meta_rows}")));
6395 }
6396
6397 let mut names: Vec<String> = Vec::new();
6399 {
6400 let mut stmt = img
6401 .prepare(&format!("PRAGMA table_info({})", quote_ident(&table.name)))
6402 .map_err(|_| invalid("image data table missing".to_owned()))?;
6403 let mut rows = stmt
6404 .query([])
6405 .map_err(|_| invalid("image data table unreadable".to_owned()))?;
6406 while let Some(row) = rows
6407 .next()
6408 .map_err(|_| invalid("image data table unreadable".to_owned()))?
6409 {
6410 names.push(
6411 row.get::<_, String>(1)
6412 .map_err(|_| invalid("image data table unreadable".to_owned()))?,
6413 );
6414 }
6415 }
6416 let mut expected: Vec<&str> = table.columns.iter().map(|c| c.name.as_str()).collect();
6417 expected.push("_syncular_version");
6418 if names.len() != expected.len() || names.iter().zip(expected.iter()).any(|(a, b)| a != b) {
6419 return Err(SectionError::Abort(
6420 "sync.schema_mismatch".to_owned(),
6421 "sqlite segment columns do not match the generated schema (§5.3)".to_owned(),
6422 ));
6423 }
6424
6425 if first_fresh_page {
6428 self.purge_scope_rows(&table.name, effective)
6429 .map_err(|()| SectionError::FailClosed)?;
6430 }
6431 self.overlay_dirty.set(true);
6437 let bulk_indexes: Vec<&crate::schema::IndexSchema> = if first_fresh_page {
6445 table.indexes.iter().filter(|i| !i.unique).collect()
6446 } else {
6447 Vec::new()
6448 };
6449 for index in &bulk_indexes {
6450 let index_name = quote_ident(&format!("_syncular_base_{}", index.name));
6451 self.conn
6452 .execute(&format!("DROP INDEX IF EXISTS {index_name}"), [])
6453 .map_err(|e| invalid(e.to_string()))?;
6454 }
6455 let insert = self.insert_row_sql(&base_table(&table.name), table);
6456 let applied = {
6457 let mut ins = self
6458 .conn
6459 .prepare_cached(&insert)
6460 .map_err(|e| invalid(e.to_string()))?;
6461 let column_list: Vec<String> = names.iter().map(|n| quote_ident(n)).collect();
6462 let mut stmt = img
6463 .prepare(&format!(
6464 "SELECT {} FROM {}",
6465 column_list.join(", "),
6466 quote_ident(&table.name)
6467 ))
6468 .map_err(|_| invalid("image data table unreadable".to_owned()))?;
6469 let mut rows = stmt
6470 .query([])
6471 .map_err(|_| invalid("image data table unreadable".to_owned()))?;
6472 let version_index = table.columns.len();
6473 let mut applied = 0u32;
6474 while let Some(row) = rows
6475 .next()
6476 .map_err(|_| invalid("image row unreadable".to_owned()))?
6477 {
6478 for (i, column) in table.columns.iter().enumerate() {
6479 let cell = row
6480 .get_ref(i)
6481 .map_err(|_| invalid("image row unreadable".to_owned()))?;
6482 let param = image_cell_param(column, cell).map_err(&invalid)?;
6483 ins.raw_bind_parameter(i + 1, param)
6484 .map_err(|e| invalid(e.to_string()))?;
6485 }
6486 let version: i64 = row
6487 .get(version_index)
6488 .map_err(|_| invalid("image row unreadable".to_owned()))?;
6489 if version < 1 {
6490 return Err(invalid(format!(
6491 "row _syncular_version must be >= 1, got {version}"
6492 )));
6493 }
6494 ins.raw_bind_parameter(version_index + 1, version)
6495 .map_err(|e| invalid(e.to_string()))?;
6496 ins.raw_execute().map_err(|e| invalid(e.to_string()))?;
6497 applied += 1;
6498 }
6499 applied
6500 };
6501 for index in &bulk_indexes {
6502 let index_name = quote_ident(&format!("_syncular_base_{}", index.name));
6503 let cols_sql = index
6504 .columns
6505 .iter()
6506 .map(|c| quote_ident(c))
6507 .collect::<Vec<_>>()
6508 .join(", ");
6509 self.conn
6510 .execute(
6511 &format!(
6512 "CREATE INDEX IF NOT EXISTS {index_name} ON {} ({cols_sql})",
6513 base_table(&table.name)
6514 ),
6515 [],
6516 )
6517 .map_err(|e| invalid(e.to_string()))?;
6518 }
6519 if i64::from(applied) != row_count {
6520 return Err(invalid(format!(
6521 "image holds {applied} rows, descriptor says {row_count}"
6522 )));
6523 }
6524 Ok(applied)
6525 }
6526
6527 fn compile_local_data_purge(
6530 &self,
6531 input: &LocalDataPurgeInput,
6532 ) -> Result<(Vec<CompiledLocalDataPurgeTarget>, String), String> {
6533 let invalid = |message: String| format!("sync.invalid_request: {message}");
6534 if input.purge_id.is_empty()
6535 || input.purge_id.len() > 128
6536 || !is_local_operation_code_like(&input.purge_id)
6537 {
6538 return Err(invalid(
6539 "local purge purgeId must be a 1–128 character code-like identifier".to_owned(),
6540 ));
6541 }
6542 if input.targets.is_empty() || input.targets.len() > MAX_LOCAL_PURGE_TARGETS {
6543 return Err(invalid(format!(
6544 "local purge needs between 1 and {MAX_LOCAL_PURGE_TARGETS} targets"
6545 )));
6546 }
6547
6548 let mut deduplicated: BTreeMap<
6549 String,
6550 (CompiledLocalDataPurgeTarget, LocalDataPurgeTarget),
6551 > = BTreeMap::new();
6552 for target in &input.targets {
6553 let table = self.schema.table(&target.table).ok_or_else(|| {
6554 invalid(format!(
6555 "local purge names unknown table {:?}",
6556 target.table
6557 ))
6558 })?;
6559 if target.selectors.is_empty() || target.selectors.len() > MAX_LOCAL_PURGE_SELECTORS {
6560 return Err(invalid(format!(
6561 "local purge target {:?} needs between 1 and {MAX_LOCAL_PURGE_SELECTORS} selectors",
6562 target.table
6563 )));
6564 }
6565 let mut selectors = Vec::with_capacity(target.selectors.len());
6566 let mut canonical_selectors = BTreeMap::new();
6567 for (column_name, raw_values) in &target.selectors {
6568 let Some((column_index, column)) = table
6569 .columns
6570 .iter()
6571 .enumerate()
6572 .find(|(_, column)| column.name == *column_name)
6573 else {
6574 return Err(invalid(format!(
6575 "local purge target {:?} names unknown column {:?}",
6576 target.table, column_name
6577 )));
6578 };
6579 let encrypted = table
6580 .encrypted_columns
6581 .iter()
6582 .any(|candidate| candidate.index == column_index);
6583 if column.ty != ColumnType::String || encrypted {
6584 return Err(invalid(format!(
6585 "local purge selector {:?}.{:?} must be a plaintext string column",
6586 target.table, column_name
6587 )));
6588 }
6589 if raw_values.is_empty() || raw_values.len() > MAX_LOCAL_PURGE_VALUES {
6590 return Err(invalid(format!(
6591 "local purge selector {:?}.{:?} needs between 1 and {MAX_LOCAL_PURGE_VALUES} values",
6592 target.table, column_name
6593 )));
6594 }
6595 let mut values = raw_values.clone();
6596 values.sort();
6597 values.dedup();
6598 if values.iter().any(|value| {
6599 value.is_empty()
6600 || value.len() > MAX_LOCAL_PURGE_VALUE_LENGTH
6601 || !is_local_operation_code_like(value)
6602 }) {
6603 return Err(invalid(format!(
6604 "local purge selector values must be 1–{MAX_LOCAL_PURGE_VALUE_LENGTH} character code-like identifiers"
6605 )));
6606 }
6607 selectors.push((column_name.clone(), values.clone()));
6608 canonical_selectors.insert(column_name.clone(), values);
6609 }
6610 let canonical = LocalDataPurgeTarget {
6611 table: target.table.clone(),
6612 selectors: canonical_selectors,
6613 };
6614 let key = serde_json::to_string(&canonical).map_err(|error| error.to_string())?;
6615 deduplicated.insert(
6616 key,
6617 (
6618 CompiledLocalDataPurgeTarget {
6619 table: target.table.clone(),
6620 selectors,
6621 },
6622 canonical,
6623 ),
6624 );
6625 }
6626 let targets = deduplicated
6627 .values()
6628 .map(|(compiled, _)| compiled.clone())
6629 .collect::<Vec<_>>();
6630 let canonical = deduplicated
6631 .values()
6632 .map(|(_, target)| target.clone())
6633 .collect::<Vec<_>>();
6634 let canonical_plan =
6635 serde_json::to_string(&canonical).map_err(|error| error.to_string())?;
6636 Ok((targets, canonical_plan))
6637 }
6638
6639 fn local_purge_base_row_ids(
6640 &self,
6641 targets: &[CompiledLocalDataPurgeTarget],
6642 ) -> Result<BTreeMap<String, BTreeSet<String>>, String> {
6643 let mut by_table: BTreeMap<String, BTreeSet<String>> = BTreeMap::new();
6644 for target in targets {
6645 let table = self
6646 .schema
6647 .table(&target.table)
6648 .ok_or_else(|| format!("sync.invalid_request: unknown table {:?}", target.table))?;
6649 let mut clauses = Vec::with_capacity(target.selectors.len());
6650 let mut params = Vec::new();
6651 for (column, values) in &target.selectors {
6652 clauses.push(format!(
6653 "{} IN ({})",
6654 quote_ident(column),
6655 values.iter().map(|_| "?").collect::<Vec<_>>().join(", ")
6656 ));
6657 params.extend(values.iter().cloned().map(SqlValue::Text));
6658 }
6659 let sql = format!(
6660 "SELECT CAST({} AS TEXT) FROM {} WHERE {}",
6661 quote_ident(&table.primary_key),
6662 base_table(&target.table),
6663 clauses.join(" AND ")
6664 );
6665 let mut statement = self.conn.prepare(&sql).map_err(|error| error.to_string())?;
6666 let rows = statement
6667 .query_map(rusqlite::params_from_iter(params), |row| {
6668 row.get::<_, String>(0)
6669 })
6670 .map_err(|error| error.to_string())?;
6671 let ids = by_table.entry(target.table.clone()).or_default();
6672 for row in rows {
6673 ids.insert(row.map_err(|error| error.to_string())?);
6674 }
6675 }
6676 Ok(by_table)
6677 }
6678
6679 fn local_purge_values_match(
6680 target: &CompiledLocalDataPurgeTarget,
6681 values: &Map<String, Value>,
6682 ) -> bool {
6683 target.selectors.iter().all(|(column, allowed)| {
6684 values
6685 .get(column)
6686 .and_then(Value::as_str)
6687 .is_some_and(|value| allowed.iter().any(|candidate| candidate == value))
6688 })
6689 }
6690
6691 pub fn purge_local_data(
6695 &mut self,
6696 input: &LocalDataPurgeInput,
6697 ) -> Result<LocalDataPurgeResult, String> {
6698 let (targets, canonical_plan) = self.compile_local_data_purge(input)?;
6699 let meta_key = format!("localPurge:{}", input.purge_id);
6700 let applied_plan = self
6701 .conn
6702 .query_row(
6703 "SELECT value FROM _syncular_meta WHERE key = ?1",
6704 rusqlite::params![meta_key],
6705 |row| row.get::<_, String>(0),
6706 )
6707 .optional()
6708 .map_err(|error| error.to_string())?;
6709 if let Some(applied_plan) = applied_plan {
6710 if applied_plan != canonical_plan {
6711 return Err(format!(
6712 "sync.invalid_request: local purge id {:?} was already used with a different plan",
6713 input.purge_id
6714 ));
6715 }
6716 return Ok(LocalDataPurgeResult {
6717 already_applied: true,
6718 purged_rows: 0,
6719 dropped_commits: 0,
6720 });
6721 }
6722
6723 let prior_outbox = self.outbox.clone();
6724 let prior_rejection_count = self.rejections.len();
6725 let prior_overlay_dirty = self.overlay_dirty.get();
6726 self.begin_observation("syncular_local_purge")?;
6727 let applied = (|| -> Result<(ChangeAccumulator, LocalDataPurgeResult), String> {
6728 let row_ids = self.local_purge_base_row_ids(&targets)?;
6729 let doomed = self
6730 .outbox
6731 .iter()
6732 .filter(|commit| {
6733 commit.ops.iter().any(|operation| {
6734 let matching_targets = targets
6735 .iter()
6736 .filter(|target| target.table == operation.table)
6737 .collect::<Vec<_>>();
6738 if matching_targets.is_empty() {
6739 return false;
6740 }
6741 if row_ids
6742 .get(&operation.table)
6743 .is_some_and(|ids| ids.contains(&operation.row_id))
6744 {
6745 return true;
6746 }
6747 operation.values.as_ref().is_some_and(|values| {
6748 matching_targets
6749 .iter()
6750 .any(|target| Self::local_purge_values_match(target, values))
6751 })
6752 })
6753 })
6754 .cloned()
6755 .collect::<Vec<_>>();
6756 let doomed_ids = doomed
6757 .iter()
6758 .map(|commit| commit.client_commit_id.clone())
6759 .collect::<BTreeSet<_>>();
6760 let mut batch = ChangeAccumulator::default();
6761 let mut rejections = Vec::new();
6762 for commit in &doomed {
6763 for operation in &commit.ops {
6764 batch.table(&operation.table);
6765 }
6766 let results = commit
6767 .ops
6768 .iter()
6769 .enumerate()
6770 .map(|(op_index, operation)| {
6771 let rejection = RejectionRecord {
6772 client_commit_id: commit.client_commit_id.clone(),
6773 op_index: op_index as i32,
6774 code: "client.local_data_purged".to_owned(),
6775 message: "the commit was dropped by an application-authorized local data purge".to_owned(),
6776 retryable: false,
6777 details: None,
6778 operation: Some(CommitOperation::from(operation)),
6779 };
6780 rejections.push(rejection.clone());
6781 CommitOperationOutcome::Error { rejection }
6782 })
6783 .collect::<Vec<_>>();
6784 self.persist_commit_outcome(
6785 &commit.client_commit_id,
6786 CommitOutcomeStatus::Rejected,
6787 &results,
6788 Some(&commit.ops),
6789 )?;
6790 self.delete_outbox_persisted(&commit.client_commit_id)?;
6791 }
6792 if !doomed.is_empty() {
6793 self.outbox
6794 .retain(|commit| !doomed_ids.contains(&commit.client_commit_id));
6795 self.rejections.extend(rejections);
6796 self.prune_commit_outcomes()?;
6797 self.overlay_dirty.set(true);
6798 batch.status = true;
6799 batch.rejections = true;
6800 batch.outcomes = true;
6801 }
6802
6803 let mut purged_rows = 0usize;
6804 for (table, ids) in &row_ids {
6805 if ids.is_empty() {
6806 continue;
6807 }
6808 batch.table(table);
6809 for row_id in ids {
6810 self.delete_base_row(table, row_id)?;
6811 purged_rows += 1;
6812 }
6813 }
6814 self.rebuild_overlay_if_dirty();
6815 self.reconcile_blob_refcounts(true);
6816 self.conn
6817 .execute(
6818 "INSERT INTO _syncular_meta(key, value) VALUES (?1, ?2)",
6819 rusqlite::params![meta_key, canonical_plan],
6820 )
6821 .map_err(|error| error.to_string())?;
6822 Ok((
6823 batch,
6824 LocalDataPurgeResult {
6825 already_applied: false,
6826 purged_rows,
6827 dropped_commits: doomed.len(),
6828 },
6829 ))
6830 })();
6831 let (batch, result) = match applied {
6832 Ok(value) => value,
6833 Err(error) => {
6834 self.rollback_observation("syncular_local_purge");
6835 self.outbox = prior_outbox;
6836 self.rejections.truncate(prior_rejection_count);
6837 self.overlay_dirty.set(prior_overlay_dirty);
6838 return Err(error);
6839 }
6840 };
6841 if let Err(error) = self.finish_observation("syncular_local_purge", batch) {
6842 self.rollback_observation("syncular_local_purge");
6843 self.outbox = prior_outbox;
6844 self.rejections.truncate(prior_rejection_count);
6845 self.overlay_dirty.set(prior_overlay_dirty);
6846 return Err(error);
6847 }
6848 Ok(result)
6849 }
6850
6851 pub fn rebootstrap_local_data(
6857 &mut self,
6858 input: &LocalDataRebootstrapInput,
6859 ) -> Result<LocalDataRebootstrapResult, String> {
6860 if input.rebootstrap_id.is_empty()
6861 || input.rebootstrap_id.len() > 128
6862 || !is_local_operation_code_like(&input.rebootstrap_id)
6863 {
6864 return Err(
6865 "sync.invalid_request: local rebootstrap rebootstrapId must be a 1–128 character code-like identifier"
6866 .to_owned(),
6867 );
6868 }
6869 let meta_key = format!("localRebootstrap:{}", input.rebootstrap_id);
6870 if let Some(persisted) = self.get_meta_strict(&meta_key)? {
6871 let (retained_commits, reset_subscriptions) =
6872 decode_local_rebootstrap_receipt(&persisted)?;
6873 return Ok(LocalDataRebootstrapResult {
6874 already_applied: true,
6875 retained_commits,
6876 reset_subscriptions,
6877 });
6878 }
6879 if self.stopped || self.schema_floor.is_some() {
6880 return Err(
6881 "sync.invalid_request: local rebootstrap cannot bypass an active schema-floor stop; update the application first"
6882 .to_owned(),
6883 );
6884 }
6885
6886 let retained_commits = self.outbox.len();
6887 let reset_subscriptions = self.subs.len();
6888 let prior_subs = self.subs.clone();
6889 let prior_upgrading = self.upgrading;
6890 let prior_stopped = self.stopped;
6891 let prior_schema_floor = self.schema_floor.clone();
6892 let prior_overlay_dirty = self.overlay_dirty.get();
6893 let prior_sync_needed = self.sync_needed;
6894 let prior_sync_intents = self.sync_intent_queue.clone();
6895 let receipt = encode_local_rebootstrap_receipt(retained_commits, reset_subscriptions)?;
6896
6897 self.begin_observation("syncular_local_rebootstrap")?;
6898 let mut batch = ChangeAccumulator::default();
6899 let applied = (|| -> Result<(), String> {
6900 self.run_schema_reset_observed(&mut batch, false)?;
6901 self.conn
6902 .execute(
6903 "INSERT INTO _syncular_meta(key, value) VALUES (?1, ?2)",
6904 rusqlite::params![meta_key, receipt],
6905 )
6906 .map_err(|error| error.to_string())?;
6907 self.sync_needed = true;
6908 self.sync_intent_queue.push_back(SyncIntent::Interactive);
6909 batch.status = true;
6910 Ok(())
6911 })();
6912 if let Err(error) = applied {
6913 self.rollback_observation("syncular_local_rebootstrap");
6914 self.subs = prior_subs;
6915 self.upgrading = prior_upgrading;
6916 self.stopped = prior_stopped;
6917 self.schema_floor = prior_schema_floor;
6918 self.overlay_dirty.set(prior_overlay_dirty);
6919 self.sync_needed = prior_sync_needed;
6920 self.sync_intent_queue = prior_sync_intents;
6921 return Err(error);
6922 }
6923 if let Err(error) = self.finish_observation("syncular_local_rebootstrap", batch) {
6924 self.rollback_observation("syncular_local_rebootstrap");
6925 self.subs = prior_subs;
6926 self.upgrading = prior_upgrading;
6927 self.stopped = prior_stopped;
6928 self.schema_floor = prior_schema_floor;
6929 self.overlay_dirty.set(prior_overlay_dirty);
6930 self.sync_needed = prior_sync_needed;
6931 self.sync_intent_queue = prior_sync_intents;
6932 return Err(error);
6933 }
6934
6935 Ok(LocalDataRebootstrapResult {
6936 already_applied: false,
6937 retained_commits,
6938 reset_subscriptions,
6939 })
6940 }
6941
6942 fn purge_scope_rows(
6947 &mut self,
6948 table_name: &str,
6949 effective: &[(String, Vec<String>)],
6950 ) -> Result<(), ()> {
6951 if effective.is_empty() {
6952 return Ok(());
6953 }
6954 let table = self.schema.table(table_name).ok_or(())?.clone();
6955 let mut clauses = Vec::new();
6956 let mut params: Vec<SqlValue> = Vec::new();
6957 for (variable, values) in effective {
6958 let column = table.scope_column(variable).ok_or(())?;
6959 let placeholders: Vec<String> = values
6960 .iter()
6961 .map(|v| {
6962 params.push(SqlValue::Text(v.clone()));
6963 "?".to_owned()
6964 })
6965 .collect();
6966 clauses.push(format!(
6967 "{} IN ({})",
6968 quote_ident(column),
6969 placeholders.join(", ")
6970 ));
6971 }
6972 let sql = format!(
6973 "DELETE FROM {} WHERE {}",
6974 base_table(table_name),
6975 clauses.join(" AND ")
6976 );
6977 self.overlay_dirty.set(true);
6978 self.conn
6979 .execute(&sql, rusqlite::params_from_iter(params))
6980 .map_err(|_| ())?;
6981 Ok(())
6982 }
6983
6984 fn drop_doomed_outbox(
6987 &mut self,
6988 table_name: &str,
6989 effective: &[(String, Vec<String>)],
6990 ) -> Result<bool, String> {
6991 if effective.is_empty() {
6992 return Ok(false);
6993 }
6994 let Some(table) = self.schema.table(table_name).cloned() else {
6995 return Ok(false);
6996 };
6997 let mut mappings: Vec<(&str, &Vec<String>)> = Vec::new();
6998 for (variable, values) in effective {
6999 match table.scope_column(variable) {
7000 Some(column) => mappings.push((column, values)),
7001 None => return Ok(false), }
7003 }
7004 let doomed: Vec<OutboxCommit> = self
7005 .outbox
7006 .iter()
7007 .filter(|commit| {
7008 commit.ops.iter().any(|op| {
7009 op.upsert
7010 && op.table == table_name
7011 && op.values.as_ref().is_some_and(|values| {
7012 mappings
7013 .iter()
7014 .all(|(column, allowed)| match values.get(*column) {
7015 Some(Value::String(s)) => allowed.contains(s),
7016 Some(Value::Number(n)) => allowed.contains(&n.to_string()),
7017 _ => false,
7018 })
7019 })
7020 })
7021 })
7022 .cloned()
7023 .collect();
7024 if doomed.is_empty() {
7025 return Ok(false);
7026 }
7027 let mut rejections = Vec::new();
7028 for commit in &doomed {
7029 let results = commit
7030 .ops
7031 .iter()
7032 .enumerate()
7033 .map(|(op_index, operation)| {
7034 let rejection = RejectionRecord {
7035 client_commit_id: commit.client_commit_id.clone(),
7036 op_index: op_index as i32,
7037 code: "sync.scope_revoked".to_owned(),
7038 message: "the commit was dropped because its effective scope was revoked"
7039 .to_owned(),
7040 retryable: false,
7041 details: None,
7042 operation: Some(CommitOperation::from(operation)),
7043 };
7044 rejections.push(rejection.clone());
7045 CommitOperationOutcome::Error { rejection }
7046 })
7047 .collect::<Vec<_>>();
7048 self.persist_commit_outcome(
7049 &commit.client_commit_id,
7050 CommitOutcomeStatus::Rejected,
7051 &results,
7052 Some(&commit.ops),
7053 )?;
7054 self.delete_outbox_persisted(&commit.client_commit_id)?;
7055 }
7056 self.prune_commit_outcomes()?;
7057 let doomed_ids = doomed
7058 .iter()
7059 .map(|commit| commit.client_commit_id.as_str())
7060 .collect::<BTreeSet<_>>();
7061 self.outbox
7062 .retain(|commit| !doomed_ids.contains(commit.client_commit_id.as_str()));
7063 self.rejections.extend(rejections);
7064 self.overlay_dirty.set(true);
7065 Ok(true)
7066 }
7067
7068 pub fn upload_blob(
7075 &mut self,
7076 bytes: &[u8],
7077 media_type: Option<String>,
7078 name: Option<String>,
7079 ) -> Result<Value, String> {
7080 let blob_id = blob_id_for(bytes);
7081 let now = self.clock_now_ms();
7082 self.conn
7083 .execute(
7084 "INSERT INTO _syncular_blobs(blob_id, bytes, byte_length, media_type, refcount, created_at_ms, last_used_ms) VALUES (?,?,?,?,0,?,?)
7085 ON CONFLICT(blob_id) DO UPDATE SET last_used_ms = excluded.last_used_ms",
7086 rusqlite::params![blob_id, bytes, bytes.len() as i64, media_type, now, now],
7087 )
7088 .map_err(|e| e.to_string())?;
7089 self.conn
7090 .execute(
7091 "INSERT OR IGNORE INTO _syncular_blob_uploads(blob_id, media_type, created_at_ms) VALUES (?,?,?)",
7092 rusqlite::params![blob_id, media_type, now],
7093 )
7094 .map_err(|e| e.to_string())?;
7095 self.enforce_blob_cache_cap();
7098 let mut obj = Map::new();
7099 obj.insert("blobId".to_owned(), Value::from(blob_id));
7100 obj.insert("byteLength".to_owned(), Value::from(bytes.len() as i64));
7101 if let Some(mt) = media_type {
7102 obj.insert("mediaType".to_owned(), Value::from(mt));
7103 }
7104 if let Some(n) = name {
7105 obj.insert("name".to_owned(), Value::from(n));
7106 }
7107 Ok(Value::Object(obj))
7108 }
7109
7110 pub fn fetch_blob(
7114 &mut self,
7115 transport: &mut dyn Transport,
7116 blob_id_or_ref: &str,
7117 ) -> Result<Value, (String, String)> {
7118 let simple = |m: String| ("client.failed".to_owned(), m);
7119 let blob_id = if blob_id_or_ref.starts_with("sha256:") {
7120 blob_id_or_ref.to_owned()
7121 } else {
7122 let value: Value = serde_json::from_str(blob_id_or_ref)
7123 .map_err(|_| simple("blob ref is not JSON".to_owned()))?;
7124 value
7125 .get("blobId")
7126 .and_then(Value::as_str)
7127 .ok_or_else(|| simple("blob ref has no blobId".to_owned()))?
7128 .to_owned()
7129 };
7130 if let Some(cached) = self.get_cached_blob(&blob_id).map_err(simple)? {
7131 return Ok(cached);
7132 }
7133 let bytes = match transport
7139 .blob_download(&blob_id)
7140 .map_err(|e| (e.code, e.message))?
7141 {
7142 BlobDownload::Bytes(bytes) => bytes,
7143 BlobDownload::Url {
7144 url,
7145 url_expires_at_ms,
7146 } => {
7147 if url_expires_at_ms.is_some_and(|exp| exp <= self.clock_now_ms()) {
7149 return Err((
7150 "sync.segment_expired".to_owned(),
7151 format!(
7152 "blob url for {blob_id} expired before fetch — re-request mints a fresh url (§5.9.5)"
7153 ),
7154 ));
7155 }
7156 transport
7157 .fetch_blob_url(&url)
7158 .map_err(|e| (e.code, e.message))?
7159 }
7160 };
7161 if blob_id_for(&bytes) != blob_id {
7163 return Err(simple(format!(
7164 "blob content address mismatch for {blob_id}"
7165 )));
7166 }
7167 let now = self.clock_now_ms();
7168 self.conn
7169 .execute(
7170 "INSERT OR IGNORE INTO _syncular_blobs(blob_id, bytes, byte_length, media_type, refcount, created_at_ms, last_used_ms) VALUES (?,?,?,NULL,0,?,?)",
7171 rusqlite::params![blob_id, bytes, bytes.len() as i64, now, now],
7172 )
7173 .map_err(|e| simple(e.to_string()))?;
7174 self.enforce_blob_cache_cap();
7175 self.get_cached_blob(&blob_id)
7176 .map_err(simple)?
7177 .ok_or_else(|| simple("blob cache write failed".to_owned()))
7178 }
7179
7180 fn get_cached_blob(&self, blob_id: &str) -> Result<Option<Value>, String> {
7181 let _ = self.conn.execute(
7184 "UPDATE _syncular_blobs SET last_used_ms = ? WHERE blob_id = ?",
7185 rusqlite::params![self.clock_now_ms(), blob_id],
7186 );
7187 let mut stmt = self
7188 .conn
7189 .prepare("SELECT bytes, byte_length, media_type FROM _syncular_blobs WHERE blob_id = ?")
7190 .map_err(|e| e.to_string())?;
7191 let mut rows = stmt
7192 .query(rusqlite::params![blob_id])
7193 .map_err(|e| e.to_string())?;
7194 if let Some(row) = rows.next().map_err(|e| e.to_string())? {
7195 let bytes: Vec<u8> = row.get(0).map_err(|e| e.to_string())?;
7196 let byte_length: i64 = row.get(1).map_err(|e| e.to_string())?;
7197 let media_type: Option<String> = row.get(2).map_err(|e| e.to_string())?;
7198 let mut obj = Map::new();
7199 obj.insert("blobId".to_owned(), Value::from(blob_id.to_owned()));
7200 obj.insert("byteLength".to_owned(), Value::from(byte_length));
7201 let mut bytes_obj = Map::new();
7202 bytes_obj.insert("$bytes".to_owned(), Value::from(bytes_to_hex(&bytes)));
7203 obj.insert("bytes".to_owned(), Value::Object(bytes_obj));
7204 if let Some(mt) = media_type {
7205 obj.insert("mediaType".to_owned(), Value::from(mt));
7206 }
7207 return Ok(Some(Value::Object(obj)));
7208 }
7209 Ok(None)
7210 }
7211
7212 fn flush_blob_uploads(&mut self, transport: &mut dyn Transport) -> Result<(), TransportError> {
7214 let pending: Vec<(String, Option<String>)> = {
7215 let mut stmt = self
7216 .conn
7217 .prepare(
7218 "SELECT blob_id, media_type FROM _syncular_blob_uploads ORDER BY created_at_ms",
7219 )
7220 .map_err(|e| TransportError::new("client.failed", e.to_string()))?;
7221 let rows = stmt
7222 .query_map([], |row| {
7223 Ok((row.get::<_, String>(0)?, row.get::<_, Option<String>>(1)?))
7224 })
7225 .map_err(|e| TransportError::new("client.failed", e.to_string()))?;
7226 rows.filter_map(Result::ok).collect()
7227 };
7228 for (blob_id, media_type) in pending {
7229 let bytes: Option<Vec<u8>> = self
7230 .conn
7231 .query_row(
7232 "SELECT bytes FROM _syncular_blobs WHERE blob_id = ?",
7233 rusqlite::params![blob_id],
7234 |row| row.get(0),
7235 )
7236 .ok();
7237 if let Some(bytes) = bytes {
7238 self.upload_one(transport, &blob_id, &bytes, media_type.as_deref())?;
7239 }
7240 let _ = self.conn.execute(
7241 "DELETE FROM _syncular_blob_uploads WHERE blob_id = ?",
7242 rusqlite::params![blob_id],
7243 );
7244 }
7245 Ok(())
7246 }
7247
7248 fn upload_one(
7255 &self,
7256 transport: &mut dyn Transport,
7257 blob_id: &str,
7258 bytes: &[u8],
7259 media_type: Option<&str>,
7260 ) -> Result<(), TransportError> {
7261 match transport.blob_upload_grant(blob_id, bytes.len() as u64, media_type)? {
7262 BlobUploadGrant::Present => return Ok(()), BlobUploadGrant::Url {
7264 url,
7265 url_expires_at_ms,
7266 } => {
7267 let live = url_expires_at_ms.is_none_or(|exp| exp > self.clock_now_ms());
7268 if live && transport.blob_put_url(&url, bytes, media_type).is_ok() {
7269 return Ok(());
7270 }
7271 }
7273 BlobUploadGrant::None => {
7274 }
7276 }
7277 transport.blob_upload(blob_id, bytes, media_type)
7278 }
7279
7280 fn enforce_blob_cache_cap(&self) {
7289 let Some(max_bytes) = self.limits.blob_cache_max_bytes else {
7290 return;
7291 };
7292 let mut total: i64 = self
7293 .conn
7294 .query_row(
7295 "SELECT COALESCE(SUM(byte_length), 0) FROM _syncular_blobs",
7296 [],
7297 |row| row.get(0),
7298 )
7299 .unwrap_or(0);
7300 if total <= max_bytes {
7301 return;
7302 }
7303 let candidates: Vec<(String, i64)> = {
7304 let Ok(mut stmt) = self.conn.prepare(
7305 "SELECT blob_id, byte_length FROM _syncular_blobs
7306 WHERE refcount = 0
7307 AND blob_id NOT IN (SELECT blob_id FROM _syncular_blob_uploads)
7308 ORDER BY last_used_ms ASC, created_at_ms ASC",
7309 ) else {
7310 return;
7311 };
7312 let Ok(rows) = stmt.query_map([], |row| {
7313 Ok((row.get::<_, String>(0)?, row.get::<_, i64>(1)?))
7314 }) else {
7315 return;
7316 };
7317 rows.filter_map(Result::ok).collect()
7318 };
7319 for (blob_id, byte_length) in candidates {
7320 if total <= max_bytes {
7321 break;
7322 }
7323 let _ = self.conn.execute(
7324 "DELETE FROM _syncular_blobs WHERE blob_id = ?",
7325 rusqlite::params![blob_id],
7326 );
7327 total -= byte_length;
7328 }
7329 }
7330
7331 fn reconcile_blob_refcounts(&mut self, delete_orphans: bool) {
7335 if !self.schema_has_blobs() {
7336 return;
7337 }
7338 let mut counts: std::collections::HashMap<String, i64> = std::collections::HashMap::new();
7339 for table in self.schema.tables.clone() {
7340 let blob_cols: Vec<String> = table
7341 .columns
7342 .iter()
7343 .filter(|c| c.ty == ColumnType::BlobRef)
7344 .map(|c| c.name.clone())
7345 .collect();
7346 for column in blob_cols {
7347 let sql = format!(
7348 "SELECT {} FROM {} WHERE {} IS NOT NULL",
7349 quote_ident(&column),
7350 base_table(&table.name),
7351 quote_ident(&column)
7352 );
7353 let Ok(mut stmt) = self.conn.prepare(&sql) else {
7354 continue;
7355 };
7356 let Ok(rows) = stmt.query_map([], |row| row.get::<_, Option<String>>(0)) else {
7357 continue;
7358 };
7359 for raw in rows.flatten().flatten() {
7360 if let Ok(value) = serde_json::from_str::<Value>(&raw) {
7361 if let Some(id) = value.get("blobId").and_then(Value::as_str) {
7362 *counts.entry(id.to_owned()).or_insert(0) += 1;
7363 }
7364 }
7365 }
7366 }
7367 }
7368 let _ = self
7369 .conn
7370 .execute("UPDATE _syncular_blobs SET refcount = 0", []);
7371 for (blob_id, count) in &counts {
7372 let _ = self.conn.execute(
7373 "UPDATE _syncular_blobs SET refcount = ? WHERE blob_id = ?",
7374 rusqlite::params![count, blob_id],
7375 );
7376 }
7377 if delete_orphans {
7378 let _ = self.conn.execute(
7379 "DELETE FROM _syncular_blobs WHERE refcount = 0 AND blob_id NOT IN (SELECT blob_id FROM _syncular_blob_uploads)",
7380 [],
7381 );
7382 }
7383 }
7384
7385 fn insert_row_sql(&self, full_table: &str, table: &crate::schema::TableSchema) -> String {
7390 if let Some(sql) = self.insert_sql.borrow().get(full_table) {
7391 return sql.clone();
7392 }
7393 let mut columns: Vec<String> = table.columns.iter().map(|c| quote_ident(&c.name)).collect();
7394 columns.push(quote_ident("_syncular_version"));
7395 let placeholders: Vec<&str> = columns.iter().map(|_| "?").collect();
7396 let primary_key = quote_ident(&table.primary_key);
7397 let updates = columns
7398 .iter()
7399 .filter(|column| **column != primary_key)
7400 .map(|column| format!("{column}=excluded.{column}"))
7401 .collect::<Vec<_>>()
7402 .join(", ");
7403 let sql = format!(
7404 "INSERT INTO {full_table} ({}) VALUES ({}) ON CONFLICT ({primary_key}) DO UPDATE SET {updates}",
7405 columns.join(", "),
7406 placeholders.join(", ")
7407 );
7408 self.insert_sql
7409 .borrow_mut()
7410 .insert(full_table.to_owned(), sql.clone());
7411 sql
7412 }
7413
7414 fn write_base_row(&self, table_name: &str, row: &Row, version: i64) -> Result<(), String> {
7415 self.overlay_dirty.set(true);
7416 self.write_row(&base_table(table_name), table_name, row, version)
7417 }
7418
7419 fn write_row(
7420 &self,
7421 full_table: &str,
7422 table_name: &str,
7423 row: &Row,
7424 version: i64,
7425 ) -> Result<(), String> {
7426 let table = self
7427 .schema
7428 .table(table_name)
7429 .ok_or_else(|| format!("unknown table {table_name:?}"))?;
7430 let sql = self.insert_row_sql(full_table, table);
7431 let mut stmt = self.conn.prepare_cached(&sql).map_err(|e| e.to_string())?;
7432 let params = row
7433 .iter()
7434 .map(RowParam::Cell)
7435 .chain(std::iter::once(RowParam::Version(version)));
7436 stmt.execute(rusqlite::params_from_iter(params))
7437 .map_err(|e| e.to_string())?;
7438 Ok(())
7439 }
7440
7441 fn delete_base_row(&self, table_name: &str, row_id: &str) -> Result<(), String> {
7442 let table = self
7443 .schema
7444 .table(table_name)
7445 .ok_or_else(|| format!("unknown table {table_name:?}"))?;
7446 self.overlay_dirty.set(true);
7447 let sql = format!(
7448 "DELETE FROM {} WHERE CAST({} AS TEXT) = ?1",
7449 base_table(table_name),
7450 quote_ident(&table.primary_key)
7451 );
7452 let mut stmt = self.conn.prepare_cached(&sql).map_err(|e| e.to_string())?;
7453 stmt.execute(rusqlite::params![row_id])
7454 .map_err(|e| e.to_string())?;
7455 Ok(())
7456 }
7457
7458 fn rebuild_overlay_if_dirty(&mut self) {
7461 if self.overlay_dirty.get() {
7462 self.rebuild_overlay();
7463 }
7464 }
7465
7466 fn rebuild_overlay(&mut self) {
7470 #[cfg(test)]
7471 self.overlay_rebuild_count
7472 .set(self.overlay_rebuild_count.get() + 1);
7473 self.exec("SAVEPOINT syncular_overlay");
7474 for table in self.schema.tables.clone() {
7475 for index in &table.fts_indexes {
7476 let _ = self.drop_fts_triggers(index);
7477 }
7478 let visible = visible_table(&table.name);
7479 let base = base_table(&table.name);
7480 self.exec(&format!("DELETE FROM {visible}"));
7481 self.exec(&format!("INSERT INTO {visible} SELECT * FROM {base}"));
7482 }
7483 for commit in self.outbox.clone() {
7484 for op in &commit.ops {
7485 let Some(table) = self.schema.table(&op.table).cloned() else {
7486 continue;
7487 };
7488 if op.upsert {
7489 let Some(values) = op.values.as_ref() else {
7490 continue;
7491 };
7492 let mut row: Row = Vec::with_capacity(table.columns.len());
7493 let mut ok = true;
7494 for column in &table.columns {
7495 match json_to_column_value(column, values.get(&column.name)) {
7496 Ok(v) => row.push(v),
7497 Err(_) => {
7498 ok = false;
7499 break;
7500 }
7501 }
7502 }
7503 if ok {
7504 let _ = self.write_row(&visible_table(&table.name), &table.name, &row, -1);
7505 }
7506 } else {
7507 let sql = format!(
7508 "DELETE FROM {} WHERE CAST({} AS TEXT) = ?1",
7509 visible_table(&table.name),
7510 quote_ident(&table.primary_key)
7511 );
7512 let _ = self.conn.execute(&sql, rusqlite::params![op.row_id]);
7513 }
7514 }
7515 }
7516 for table in self.schema.tables.clone() {
7517 for index in &table.fts_indexes {
7518 let _ = self.rebuild_fts_projection(&table, index);
7519 let _ = self.create_fts_triggers(&table, index);
7520 }
7521 }
7522 self.exec("RELEASE syncular_overlay");
7523 self.overlay_dirty.set(false);
7524 }
7525
7526 fn exec(&self, sql: &str) {
7527 let _ = self.conn.execute_batch(sql);
7528 }
7529
7530 pub fn connect_realtime(&mut self, transport: &mut dyn Transport) -> Result<(), String> {
7533 if self.realtime_connected {
7534 return Ok(());
7535 }
7536 transport
7537 .realtime_connect_for_client(&self.client_id)
7538 .map_err(|e| format!("{}: {}", e.code, e.message))?;
7539 self.realtime_connected = true;
7540 Ok(())
7541 }
7542
7543 pub fn disconnect_realtime(&mut self, transport: &mut dyn Transport) {
7544 if !self.realtime_connected {
7545 return;
7546 }
7547 let _ = transport.realtime_close();
7548 self.realtime_connected = false;
7549 self.presence.clear(); }
7551
7552 pub fn set_presence(
7558 &mut self,
7559 transport: &mut dyn Transport,
7560 scope_key: &str,
7561 doc: Option<&Value>,
7562 ) -> Result<(), String> {
7563 if !self.realtime_connected {
7564 return Err("setPresence requires a connected realtime socket (§8.6)".to_string());
7565 }
7566 let text = encode_presence_publish(scope_key, doc);
7567 transport
7568 .realtime_send(&text)
7569 .map_err(|e| format!("{}: {}", e.code, e.message))
7570 }
7571
7572 pub fn presence(&self, scope_key: &str) -> Vec<PresencePeer> {
7574 self.presence
7575 .get(scope_key)
7576 .map(|peers| peers.values().cloned().collect())
7577 .unwrap_or_default()
7578 }
7579
7580 fn apply_presence(
7582 &mut self,
7583 scope_key: String,
7584 kind: Option<PresenceKind>,
7585 actor_id: Option<String>,
7586 client_id: Option<String>,
7587 doc: Option<Value>,
7588 error: Option<String>,
7589 ) {
7590 if error.is_some() {
7593 return;
7594 }
7595 let (Some(kind), Some(actor_id), Some(client_id)) = (kind, actor_id, client_id) else {
7596 return;
7597 };
7598 let peer_key = format!("{actor_id} {client_id}");
7599 match kind {
7600 PresenceKind::Leave => {
7601 if let Some(peers) = self.presence.get_mut(&scope_key) {
7602 peers.remove(&peer_key);
7603 if peers.is_empty() {
7604 self.presence.remove(&scope_key);
7605 }
7606 }
7607 }
7608 _ => {
7609 let doc = match doc {
7610 Some(Value::Object(_)) => doc.unwrap(),
7611 _ => return,
7612 };
7613 self.presence.entry(scope_key).or_default().insert(
7614 peer_key,
7615 PresencePeer {
7616 actor_id,
7617 client_id,
7618 doc,
7619 },
7620 );
7621 }
7622 }
7623 }
7624
7625 pub fn on_realtime_text(&mut self, text: &str) {
7627 match parse_control(text) {
7628 Ok(ControlMessage::Hello { requires_sync, .. }) => {
7629 if requires_sync {
7630 self.set_sync_needed(true, true);
7632 }
7633 }
7634 Ok(ControlMessage::Presence {
7635 scope_key,
7636 kind,
7637 actor_id,
7638 client_id,
7639 doc,
7640 error,
7641 ..
7642 }) => {
7643 self.apply_presence(scope_key, kind, actor_id, client_id, doc, error);
7644 }
7645 Ok(ControlMessage::Wake { .. }) => {
7646 self.set_sync_needed(true, true);
7648 }
7649 _ => {}
7650 }
7651 }
7652
7653 pub fn on_realtime_binary(&mut self, transport: &mut dyn Transport, bytes: &[u8]) {
7656 if self.stopped {
7657 return;
7658 }
7659 let message = match decode_message(bytes) {
7660 Ok(m) if m.msg_kind == MsgKind::Response => m,
7661 _ => {
7662 self.set_sync_needed(true, true);
7663 return;
7664 }
7665 };
7666 let mut frames = message.frames.into_iter();
7667 let mut applied_cursor: Option<i64> = None;
7668 let mut any_covered = false;
7669 let mut dropped = false;
7670 while let Some(frame) = frames.next() {
7671 let Frame::SubStart {
7672 id,
7673 status,
7674 effective_scopes,
7675 ..
7676 } = frame
7677 else {
7678 continue;
7679 };
7680 let mut body = Vec::new();
7681 let mut next_cursor: Option<i64> = None;
7682 for inner in frames.by_ref() {
7683 match inner {
7684 Frame::SubEnd {
7685 next_cursor: nc, ..
7686 } => {
7687 next_cursor = Some(nc);
7688 break;
7689 }
7690 Frame::Unknown { .. } => {}
7691 other => body.push(other),
7692 }
7693 }
7694 let Some(next_cursor) = next_cursor else {
7695 dropped = true;
7696 break;
7697 };
7698 let Some(sub_index) = self.subs.iter().position(|s| s.id == id) else {
7699 dropped = true;
7700 continue;
7701 };
7702 let sub = &self.subs[sub_index];
7703 if status != SubStatus::Active
7706 || sub.state != SubState::Active
7707 || sub.bootstrap_state.is_some()
7708 || !sub.synced_once
7709 {
7710 dropped = true;
7711 continue;
7712 }
7713 if next_cursor <= sub.cursor {
7714 any_covered = true;
7716 continue;
7717 }
7718 let previous_effective = self.subs[sub_index].effective.clone();
7719 let previous_cursor = self.subs[sub_index].cursor;
7720 if self.begin_observation("syncular_delta").is_err() {
7721 dropped = true;
7722 continue;
7723 }
7724 self.subs[sub_index].effective = Some(effective_scopes);
7725 let mut batch = ChangeAccumulator::default();
7726 let mut failed = false;
7727 for inner in body {
7728 if let Frame::Commit {
7729 tables, changes, ..
7730 } = inner
7731 {
7732 self.record_commit_changes(&mut batch, &tables, &changes);
7733 if self.apply_commit_changes(&tables, &changes).is_err() {
7734 failed = true;
7735 break;
7736 }
7737 }
7738 }
7739 if failed {
7740 self.rollback_observation("syncular_delta");
7741 self.subs[sub_index].effective = previous_effective;
7742 self.subs[sub_index].cursor = previous_cursor;
7743 self.overlay_dirty.set(true);
7744 self.rebuild_overlay();
7745 dropped = true;
7746 continue;
7747 }
7748 let sub = &mut self.subs[sub_index];
7749 sub.cursor = next_cursor;
7750 self.persist_sub(&self.subs[sub_index].clone());
7751 self.rebuild_overlay_if_dirty();
7752 if self.finish_observation("syncular_delta", batch).is_err() {
7753 self.rollback_observation("syncular_delta");
7754 self.subs[sub_index].effective = previous_effective;
7755 self.subs[sub_index].cursor = previous_cursor;
7756 self.overlay_dirty.set(true);
7757 self.rebuild_overlay();
7758 dropped = true;
7759 continue;
7760 }
7761 applied_cursor = Some(applied_cursor.map_or(next_cursor, |c| c.max(next_cursor)));
7762 }
7763 if let Some(cursor) = applied_cursor {
7764 self.reconcile_blob_refcounts(false);
7765 let ack = format!("{{\"type\":\"ack\",\"cursor\":{cursor}}}");
7767 let _ = transport.realtime_send(&ack);
7768 } else if !any_covered || dropped {
7769 self.set_sync_needed(true, true);
7771 }
7772 }
7773
7774 fn ack_after_pull(&mut self, transport: &mut dyn Transport) {
7778 if !self.realtime_connected {
7779 return;
7780 }
7781 let floor = self
7782 .subs
7783 .iter()
7784 .filter(|s| {
7785 s.state == SubState::Active
7786 && s.bootstrap_state.is_none()
7787 && s.synced_once
7788 && s.cursor >= 0
7789 })
7790 .map(|s| s.cursor)
7791 .min();
7792 if let Some(cursor) = floor {
7793 let ack = format!("{{\"type\":\"ack\",\"cursor\":{cursor}}}");
7794 let _ = transport.realtime_send(&ack);
7795 }
7796 }
7797}