#![cfg(feature = "offline-sync")]
#![allow(clippy::cast_possible_wrap)]
use autumn_web::sync::{
Change, ChangeOutcome, LwwResolver, Op, PgSyncBackend, PushRequest, SyncBackend, SyncScope,
};
use chrono::{DateTime, TimeZone, Utc};
use diesel::connection::SimpleConnection;
use diesel::prelude::*;
use diesel::sql_types::{BigInt, Text};
use testcontainers::ImageExt;
use testcontainers::runners::AsyncRunner;
use testcontainers_modules::postgres::Postgres;
const SCOPE: &str = SyncScope::GLOBAL;
const SEED_ROWS: usize = 20_000;
fn base_time() -> DateTime<Utc> {
Utc.with_ymd_and_hms(2026, 1, 1, 0, 0, 0).unwrap()
}
fn seeded_pk(i: usize) -> String {
format!("pk-{i:08}")
}
const fn seeded_collection(i: usize) -> &'static str {
match i % 100 {
0..=39 => "notes",
40..=59 => "contacts",
60..=74 => "tasks",
75..=89 => "files",
90..=96 => "tags",
_ => "archive",
}
}
fn seed_fixture(conn: &mut PgConnection) {
conn.batch_execute(&format!(
"INSERT INTO autumn_sync_rows \
(scope, collection, pk, payload, version, deleted, updated_at, device_id, created_version) \
SELECT \
'{SCOPE}', \
CASE WHEN i % 100 < 40 THEN 'notes' \
WHEN i % 100 < 60 THEN 'contacts' \
WHEN i % 100 < 75 THEN 'tasks' \
WHEN i % 100 < 90 THEN 'files' \
WHEN i % 100 < 97 THEN 'tags' \
ELSE 'archive' END, \
'pk-' || lpad(i::text, 8, '0'), \
jsonb_build_object( \
'title', repeat(md5(i::text), 1 + i % 10), \
'n', i, \
'note', CASE WHEN i % 7 = 0 THEN NULL ELSE md5((i * 7)::text) END \
), \
i, \
(i % 20 = 0), \
TIMESTAMPTZ '2025-01-01 00:00:00Z' + (i || ' seconds')::interval, \
'seed-device', \
i \
FROM generate_series(1, {SEED_ROWS}) AS i"
))
.expect("seed autumn_sync_rows");
conn.batch_execute(&format!(
"SELECT setval('autumn_sync_version_seq', {SEED_ROWS}, true)"
))
.expect("advance version sequence past seeded rows");
conn.batch_execute(
"UPDATE autumn_sync_rows SET device_id = 'seed-device-2' WHERE right(pk, 1) IN ('0', '1')",
)
.expect("create dead tuples");
conn.batch_execute("ANALYZE autumn_sync_rows, autumn_sync_applied, autumn_sync_horizons")
.expect("analyze");
}
struct Batch {
request: PushRequest,
existing_indices: Vec<usize>,
}
#[allow(clippy::arithmetic_side_effects)] #[allow(clippy::too_many_lines)] fn build_batch(device_id: &str, m: usize, offset: usize) -> Batch {
let (updates_n, inserts_n, deletes_n, dup_existing_pairs, dup_new_pairs) = match m {
50 => (30, 10, 4, 2, 1),
250 => (160, 50, 12, 10, 4),
1000 => (650, 200, 50, 40, 10),
_ => unreachable!("fixed batch sizes only"),
};
assert_eq!(
updates_n + inserts_n + deletes_n + dup_existing_pairs * 2 + dup_new_pairs * 2,
m,
"batch bucket sizes must sum to the batch size"
);
let base = base_time();
let mut changes = Vec::with_capacity(m);
let mut existing_indices = Vec::new();
let mut next_idx = offset + 1;
let mut seq = 0u32;
let mut next_change_id = || {
seq += 1;
format!("{device_id}-chg-{seq:06}")
};
let mut next_live_idx = || loop {
let idx = next_idx;
next_idx += 1;
if !idx.is_multiple_of(20) {
return idx;
}
};
for _ in 0..updates_n {
let idx = next_live_idx();
existing_indices.push(idx);
changes.push(Change {
change_id: next_change_id(),
collection: seeded_collection(idx).to_owned(),
pk: seeded_pk(idx),
op: Op::Upsert,
payload: Some(serde_json::json!({"edit": "update", "idx": idx})),
base_version: idx as i64,
updated_at: base + chrono::Duration::seconds(1_000_000 + idx as i64),
});
}
for k in 0..inserts_n {
changes.push(Change {
change_id: next_change_id(),
collection: "notes".to_owned(),
pk: format!("new-{device_id}-{k:06}"),
op: Op::Upsert,
payload: Some(serde_json::json!({"edit": "insert", "k": k})),
base_version: 0,
updated_at: base + chrono::Duration::seconds(2_000_000 + k as i64),
});
}
for _ in 0..deletes_n {
let idx = next_live_idx();
existing_indices.push(idx);
changes.push(Change {
change_id: next_change_id(),
collection: seeded_collection(idx).to_owned(),
pk: seeded_pk(idx),
op: Op::Delete,
payload: None,
base_version: idx as i64,
updated_at: base + chrono::Duration::seconds(3_000_000 + idx as i64),
});
}
for p in 0..dup_existing_pairs {
let idx = next_live_idx();
existing_indices.push(idx);
let pk = seeded_pk(idx);
let collection = seeded_collection(idx).to_owned();
changes.push(Change {
change_id: next_change_id(),
collection: collection.clone(),
pk: pk.clone(),
op: Op::Upsert,
payload: Some(serde_json::json!({"edit": "dup-existing-1", "p": p})),
base_version: idx as i64,
updated_at: base + chrono::Duration::seconds(4_000_000 + p as i64 * 10),
});
changes.push(Change {
change_id: next_change_id(),
collection,
pk,
op: Op::Upsert,
payload: Some(serde_json::json!({"edit": "dup-existing-2", "p": p})),
base_version: idx as i64,
updated_at: base + chrono::Duration::seconds(4_000_000 + p as i64 * 10 + 5),
});
}
for p in 0..dup_new_pairs {
let pk = format!("newdup-{device_id}-{p:06}");
changes.push(Change {
change_id: next_change_id(),
collection: "notes".to_owned(),
pk: pk.clone(),
op: Op::Upsert,
payload: Some(serde_json::json!({"edit": "dup-new-1", "p": p})),
base_version: 0,
updated_at: base + chrono::Duration::seconds(5_000_000 + p as i64 * 10),
});
changes.push(Change {
change_id: next_change_id(),
collection: "notes".to_owned(),
pk,
op: Op::Upsert,
payload: Some(serde_json::json!({"edit": "dup-new-2", "p": p})),
base_version: 0,
updated_at: base + chrono::Duration::seconds(5_000_000 + p as i64 * 10 + 5),
});
}
assert_eq!(changes.len(), m);
Batch {
request: PushRequest {
device_id: device_id.to_owned(),
changes,
},
existing_indices,
}
}
fn assert_clean_outcomes(request: &PushRequest, outcomes: &[ChangeOutcome]) {
assert_eq!(outcomes.len(), request.changes.len());
for (change, outcome) in request.changes.iter().zip(outcomes) {
let is_second_of_dup_pair = change
.payload
.as_ref()
.and_then(|p| p.get("edit"))
.and_then(|v| v.as_str())
.is_some_and(|edit| edit.ends_with("-2"));
match outcome {
ChangeOutcome::Applied { .. } => {
assert!(
!is_second_of_dup_pair,
"expected the second half of a dup pair to conflict (Resolved), \
got Applied for {}",
change.change_id
);
}
ChangeOutcome::Resolved { row } => {
assert!(
is_second_of_dup_pair,
"unexpected Resolved outcome for {}",
change.change_id
);
assert_eq!(
change.payload.as_ref(),
row.payload.as_ref(),
"LwwResolver must pick the later (client/second) write for {}",
change.change_id
);
}
ChangeOutcome::AlreadyApplied { .. } => {
panic!(
"unexpected AlreadyApplied on a first push: {}",
change.change_id
);
}
}
}
}
#[derive(QueryableByName, Debug)]
struct StatementRow {
#[diesel(sql_type = Text)]
query: String,
#[diesel(sql_type = BigInt)]
calls: i64,
#[diesel(sql_type = BigInt)]
buffers: i64,
}
fn reset_stats(conn: &mut PgConnection) {
conn.batch_execute("SELECT pg_stat_statements_reset()")
.expect("reset pg_stat_statements");
}
fn print_profile(conn: &mut PgConnection, label: &str) -> i64 {
println!("\n=== pg_stat_statements: {label} ===");
let by_calls = diesel::sql_query(
"SELECT query, calls, (shared_blks_hit + shared_blks_read) AS buffers \
FROM pg_stat_statements \
WHERE (query ILIKE '%autumn_sync%' OR query ILIKE '%nextval%' \
OR query ILIKE '%generate_series%') \
AND query NOT ILIKE '%pg_stat_statements%' \
ORDER BY calls DESC LIMIT 10",
)
.load::<StatementRow>(conn)
.expect("query pg_stat_statements by calls");
let mut total_calls = 0i64;
let mut total_buffers = 0i64;
println!("-- top by calls --");
for row in &by_calls {
total_calls += row.calls;
total_buffers += row.buffers;
println!(
"calls={:<6} buffers={:<6} {}",
row.calls,
row.buffers,
row.query.split_whitespace().collect::<Vec<_>>().join(" ")
);
}
let by_buffers = diesel::sql_query(
"SELECT query, calls, (shared_blks_hit + shared_blks_read) AS buffers \
FROM pg_stat_statements \
WHERE (query ILIKE '%autumn_sync%' OR query ILIKE '%nextval%' \
OR query ILIKE '%generate_series%') \
AND query NOT ILIKE '%pg_stat_statements%' \
ORDER BY buffers DESC LIMIT 10",
)
.load::<StatementRow>(conn)
.expect("query pg_stat_statements by buffers");
println!("-- top by buffers --");
for row in &by_buffers {
println!(
"buffers={:<6} calls={:<6} {}",
row.buffers,
row.calls,
row.query.split_whitespace().collect::<Vec<_>>().join(" ")
);
}
println!(
"TOTAL sync-table statement calls={total_calls} buffers={total_buffers} (this batch, all statement shapes)"
);
total_calls
}
#[derive(QueryableByName, Debug)]
struct ExplainLine {
#[diesel(sql_type = Text, column_name = "QUERY PLAN")]
line: String,
}
fn explain(conn: &mut PgConnection, label: &str, sql: &str) {
println!("\n=== EXPLAIN (ANALYZE, BUFFERS, VERBOSE, SETTINGS): {label} ===");
println!("{sql}");
let lines = diesel::sql_query(format!(
"EXPLAIN (ANALYZE, BUFFERS, VERBOSE, SETTINGS) {sql}"
))
.load::<ExplainLine>(conn)
.expect("explain");
for line in lines {
println!("{}", line.line);
}
}
#[derive(QueryableByName, Debug, PartialEq, Eq)]
struct RowDump {
#[diesel(sql_type = Text)]
collection: String,
#[diesel(sql_type = Text)]
pk: String,
#[diesel(sql_type = diesel::sql_types::Nullable<Text>)]
payload: Option<String>,
#[diesel(sql_type = BigInt)]
version: i64,
#[diesel(sql_type = diesel::sql_types::Bool)]
deleted: bool,
#[diesel(sql_type = Text)]
device_id: String,
#[diesel(sql_type = BigInt)]
created_version: i64,
}
fn dump_scope(conn: &mut PgConnection, prefix: &str) -> Vec<RowDump> {
diesel::sql_query(format!(
"SELECT collection, pk, payload::text AS payload, version, deleted, device_id, created_version \
FROM autumn_sync_rows WHERE scope = '{SCOPE}' AND pk LIKE '{prefix}%' \
ORDER BY collection, pk"
))
.load::<RowDump>(conn)
.expect("dump scope")
}
#[tokio::test]
#[ignore = "requires Docker (testcontainers)"]
#[allow(clippy::too_many_lines)] async fn offline_sync_push_batching_profile() {
let container = Postgres::default()
.with_tag("16-alpine")
.with_cmd([
"-c",
"fsync=off",
"-c",
"shared_preload_libraries=pg_stat_statements",
"-c",
"pg_stat_statements.track=all",
"-c",
"pg_stat_statements.max=2000",
])
.start()
.await
.expect("failed to start postgres container");
let host = container.get_host().await.expect("host");
let port = container.get_host_port_ipv4(5432).await.expect("port");
let url = format!("postgres://postgres:postgres@{host}:{port}/postgres");
let backend = PgSyncBackend::new(url.clone());
backend.ensure_schema().expect("ensure schema");
let mut conn = PgConnection::establish(&url).expect("db connection");
conn.batch_execute("CREATE EXTENSION IF NOT EXISTS pg_stat_statements")
.expect("create pg_stat_statements extension");
seed_fixture(&mut conn);
let resolver = LwwResolver;
let sizes_and_offsets = [(50usize, 0usize), (250, 5_000), (1000, 10_000)];
let mut calls_by_size = Vec::new();
for (size, offset) in sizes_and_offsets {
let device_id = format!("device-{size}");
let batch = build_batch(&device_id, size, offset);
reset_stats(&mut conn);
let outcomes = backend
.apply_push(SCOPE, &batch.request, &resolver)
.expect("apply_push")
.outcomes;
assert_clean_outcomes(&batch.request, &outcomes);
let total_calls = print_profile(&mut conn, &format!("push batch of {size} changes"));
calls_by_size.push((size, total_calls));
let dump = dump_scope(&mut conn, &format!("new%-{device_id}-"));
println!("-- row dump (new/newdup pks), {} rows --", dump.len());
for row in &dump {
println!("{row:?}");
}
let existing_dump: Vec<RowDump> = diesel::sql_query(format!(
"SELECT collection, pk, payload::text AS payload, version, deleted, device_id, created_version \
FROM autumn_sync_rows WHERE scope = '{SCOPE}' AND pk = ANY($1) ORDER BY pk"
))
.bind::<diesel::sql_types::Array<Text>, _>(
batch
.existing_indices
.iter()
.map(|i| seeded_pk(*i))
.collect::<Vec<_>>(),
)
.load(&mut conn)
.expect("dump existing pks");
println!(
"-- row dump (existing pks touched by this batch), {} rows --",
existing_dump.len()
);
for row in &existing_dump {
println!("{row:?}");
}
}
println!("\n=== statement-count scaling ===");
for (size, calls) in &calls_by_size {
println!("batch size={size:<5} total sync-table statement calls={calls}");
}
let (replay_size, replay_offset) = (250, 5_000);
let replay_device = format!("device-{replay_size}");
let replay_batch = build_batch(&replay_device, replay_size, replay_offset);
reset_stats(&mut conn);
let replay_outcomes = backend
.apply_push(SCOPE, &replay_batch.request, &resolver)
.expect("replay apply_push")
.outcomes;
for (change, outcome) in replay_batch.request.changes.iter().zip(&replay_outcomes) {
assert!(
matches!(
outcome,
ChangeOutcome::AlreadyApplied { .. } | ChangeOutcome::Resolved { .. }
),
"replay of {} must not apply again, got {outcome:?}",
change.change_id
);
}
print_profile(&mut conn, "replay of the 250-change batch (idempotency)");
explain(
&mut conn,
"dedup check (typical case: miss)",
&format!(
"SELECT version, resolved_row FROM autumn_sync_applied \
WHERE scope = '{SCOPE}' AND device_id = 'device-1000' AND change_id = 'device-1000-chg-000001'"
),
);
explain(
&mut conn,
"current-row fetch FOR UPDATE",
&format!(
"SELECT collection, pk, payload, version, deleted, updated_at, device_id, created_version \
FROM autumn_sync_rows WHERE scope = '{SCOPE}' AND collection = '{}' AND pk = '{}' FOR UPDATE",
seeded_collection(10_001),
seeded_pk(10_001)
),
);
}