Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 19 additions & 1 deletion API.md
Original file line number Diff line number Diff line change
Expand Up @@ -783,7 +783,17 @@ The send path streams payloads through [`cloudsync_payload_chunks()`](#cloudsync

Chunk transport is transparent to the CloudSync backend. Each chunk is sent as a normal `/apply` payload, either inline as a base64 `blob` or through the upload `url` path. There is no separate chunk flag: old payloads, monolithic payloads, and v3 fragment payloads are distinguished by the payload format itself.

**Parameters:** None.
**Parameters:**

| Parameter | Type | Description |
|---|---|---|
| `max_db_versions` | INTEGER | Optional. Sends at most this many local database versions in this call, instead of the whole backlog. Must be greater than zero. |

Without the argument the call sends every unsent change as a single all-or-nothing batch, which the server confirms only once every chunk has applied. On a large backlog — a long offline period, or a bulk import — that batch can be big enough to fail repeatedly and be re-sent whole each time. Passing `max_db_versions` splits the upload into several smaller batches, each confirmed independently, so a failure costs one bounded window instead of the entire backlog.

The argument counts local database versions, not a version number: each version is one local transaction, so `cloudsync_network_send_changes(50)` uploads at most fifty transactions. Received changes share the same version counter, so versions holding no local change are skipped rather than consuming the budget — the call always makes progress while anything is pending.

Call it repeatedly with the same value until `send.status` is no longer `"out-of-sync"`. `send.localVersion` always reports the newest local version and `send.serverVersion` what the server has confirmed, so the two can be compared to see what is left. When nothing is pending the call sends nothing and leaves the send checkpoint where it was.

**Returns:** A JSON string with the send result:

Expand All @@ -804,6 +814,14 @@ Chunk transport is transparent to the CloudSync backend. Each chunk is sent as a
SELECT cloudsync_network_send_changes();
-- '{"send":{"status":"synced","localVersion":5,"serverVersion":5,"chunks":1,"bytes":2048}}'

-- Bounded: upload at most 2 local transactions per call, repeating until synced.
SELECT cloudsync_network_send_changes(2);
-- '{"send":{"status":"out-of-sync","localVersion":5,"serverVersion":2,"chunks":1,"bytes":900}}'
SELECT cloudsync_network_send_changes(2);
-- '{"send":{"status":"out-of-sync","localVersion":5,"serverVersion":4,"chunks":1,"bytes":880}}'
SELECT cloudsync_network_send_changes(2);
-- '{"send":{"status":"synced","localVersion":5,"serverVersion":5,"chunks":1,"bytes":460}}'

-- With a server-reported failure (e.g. unknown schema hash on the server side):
-- '{"send":{"status":"out-of-sync","localVersion":1,"serverVersion":0,"chunks":1,"bytes":512,"lastFailure":{"jobId":44961,"code":"internal_error","stage":"apply_payload","message":"cloudsync operation failed: Cannot apply the received payload because the schema hash is unknown 4288148391734624266.","retryable":true,"failedAt":"2026-04-15T22:21:09.018606Z"}}}'
```
Expand Down
4 changes: 4 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,10 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/).

## [Unreleased]

### Added

- **`cloudsync_network_send_changes()` accepts an optional limit on how many local database versions to send**, so a large backlog can be uploaded in bounded steps instead of one batch. A send is all or nothing: the server confirms the window only once every chunk of the batch has applied, and a failed batch is re-sent whole. After a long offline period or a bulk import that batch can be large enough to keep failing, and each attempt re-uploads everything. `cloudsync_network_send_changes(max_db_versions)` sends at most that many local transactions, so each call is an independently confirmed batch and a failure costs one bounded window rather than the whole backlog. Call it repeatedly with the same value until `send.status` leaves `out-of-sync`; `send.localVersion` keeps reporting the newest local version so the remaining backlog stays visible. Received changes share the database version counter, so versions holding no local change are skipped instead of consuming the budget. The no-argument form is unchanged.

### Fixed

- **SQLite: paging a chunked download no longer gets slower with every chunk.** The positional cursor on `cloudsync_payload_chunks` was meant to seek straight to where the previous call stopped, but it stated its resume point only inside `(db_version > ? OR (db_version = ? AND seq >= ?))`, whose two arms carry distinct parameters. SQLite does not derive a range from that, so the scan over `cloudsync_changes` ran with an upper bound only and re-read the window from the beginning on every call, discarding rows until it reached the resume point — making a full drain quadratic in the number of chunks, and long enough on a large tenant to hit a server-side deadline and never complete. An explicit `db_version >= ?` is now stated alongside the disjunction; it selects exactly the same rows and lets the scan seek. Locally, draining a 188-chunk window went from 6110 ms to 210 ms, and the per-chunk cost no longer depends on how large the window is. PostgreSQL was never affected.
Expand Down
156 changes: 135 additions & 21 deletions src/network/network.c
Original file line number Diff line number Diff line change
Expand Up @@ -1889,10 +1889,86 @@ void cloudsync_network_has_unsent_changes (sqlite3_context *context, int argc, s
sqlite3_result_int(context, (last_optimistic_version >= 0 && last_optimistic_version < last_local_change));
}

// Window for a bounded send, in one pass over the metadata tables: the db_version of
// the max_versions-th local change after the checkpoint (or the newest one when fewer
// remain), plus the newest local db_version so the caller can report the backlog.
//
// Deliberately not a cloudsync_changes query: a site_id-only constraint puts that vtab
// on its full-scan plan and materialises every value through cloudsync_col_value(),
// which is the cost bounded sends exist to avoid. Local changes carry site_id 0 in the
// metadata tables, and "db_version > since" seeks the (db_version) index, so the work
// is proportional to the pending backlog rather than to the table.
//
// Counting local db_versions rather than taking a raw upper bound also keeps the window
// non-empty whenever anything is pending: db_version is shared with merges, so a range
// picked blindly can hold no local change at all and make no progress.
static int network_send_window (sqlite3 *db, int64_t since, int64_t max_versions,
int64_t *until_out, int64_t *max_local_out) {
*until_out = since;
*max_local_out = since;

const char *build_sql =
"SELECT group_concat('SELECT db_version FROM \"' || format('%w',tbl_name) || "
"'\" WHERE site_id=0 AND db_version>?1', ' UNION ') "
"FROM sqlite_master WHERE type='table' AND tbl_name LIKE '%_cloudsync'";
sqlite3_stmt *vm = NULL;
int rc = sqlite3_prepare_v2(db, build_sql, -1, &vm, NULL);
if (rc != SQLITE_OK) return rc;
char *unions = NULL;
if (sqlite3_step(vm) == SQLITE_ROW && sqlite3_column_type(vm, 0) != SQLITE_NULL) {
unions = sqlite3_mprintf("%s", (const char *)sqlite3_column_text(vm, 0));
}
sqlite3_finalize(vm);
if (!unions) return SQLITE_OK; // no synced tables: nothing local to send

char *sql = sqlite3_mprintf(
"WITH v(db_version) AS (%s) "
"SELECT (SELECT db_version FROM v ORDER BY db_version LIMIT 1 OFFSET ?2), (SELECT max(db_version) FROM v)",
unions);
sqlite3_free(unions);
if (!sql) return SQLITE_NOMEM;

vm = NULL;
rc = sqlite3_prepare_v2(db, sql, -1, &vm, NULL);
sqlite3_free(sql);
if (rc != SQLITE_OK) return rc;
sqlite3_bind_int64(vm, 1, since);
sqlite3_bind_int64(vm, 2, max_versions - 1);
rc = sqlite3_step(vm);
if (rc == SQLITE_ROW) {
int64_t max_local = (sqlite3_column_type(vm, 1) == SQLITE_NULL) ? since : sqlite3_column_int64(vm, 1);
// Fewer pending versions than requested: the window is the whole backlog.
int64_t until = (sqlite3_column_type(vm, 0) == SQLITE_NULL) ? max_local : sqlite3_column_int64(vm, 0);
*max_local_out = max_local;
*until_out = until;
rc = SQLITE_OK;
} else if (rc == SQLITE_DONE) {
rc = SQLITE_OK;
}
sqlite3_finalize(vm);
return rc;
}

int cloudsync_network_send_changes_internal (sqlite3_context *context, int argc, sqlite3_value **argv, sync_result *out) {
DEBUG_FUNCTION("cloudsync_network_send_changes");
UNUSED_PARAMETER(argc);
UNUSED_PARAMETER(argv);

// One optional argument caps this call at that many local db_versions, so a large
// backlog uploads as several bounded all-or-nothing batches instead of one unbounded
// batch that must succeed whole. Every chunk of a bounded batch announces the same
// window (the vtab reports watermark = until), so the batch stays coherent.
bool bounded = (argc == 1);
int64_t max_versions = 0;
if (bounded) {
if (sqlite3_value_type(argv[0]) != SQLITE_INTEGER) {
sqlite3_result_error(context, "cloudsync_network_send_changes: max_db_versions must be an integer.", -1);
return SQLITE_ERROR;
}
max_versions = sqlite3_value_int64(argv[0]);
if (max_versions <= 0) {
sqlite3_result_error(context, "cloudsync_network_send_changes: max_db_versions must be greater than zero.", -1);
return SQLITE_ERROR;
}
}

// retrieve global context
cloudsync_context *data = (cloudsync_context *)sqlite3_user_data(context);
Expand All @@ -1907,17 +1983,40 @@ int cloudsync_network_send_changes_internal (sqlite3_context *context, int argc,
}

sqlite3 *db = sqlite3_context_db_handle(context);

// Resolve the window, and the newest local db_version for the backlog status, in
// one pass. until == db_version means nothing local is pending.
int64_t until = db_version;
int64_t max_local = db_version;
if (bounded) {
int wrc = network_send_window(db, db_version, max_versions, &until, &max_local);
if (wrc != SQLITE_OK) {
sqlite3_result_error(context, sqlite3_errmsg(db), -1);
sqlite3_result_error_code(context, wrc);
return wrc;
}
}
// cloudsync_payload_chunks reads until_db_version = 0 as "no upper bound", so an
// empty window has to be recognised here instead of being passed down.
bool empty_window = bounded && (until == db_version);

sqlite3_stmt *stmt = NULL;
const char *chunk_sql =
"SELECT payload, payload_size, watermark_db_version, is_final "
"FROM cloudsync_payload_chunks WHERE since_db_version = ?";
int rc = sqlite3_prepare_v2(db, chunk_sql, -1, &stmt, NULL);
if (rc != SQLITE_OK) {
sqlite3_result_error(context, sqlite3_errmsg(db), -1);
sqlite3_result_error_code(context, rc);
return rc;
int rc = SQLITE_OK;
if (!empty_window) {
const char *chunk_sql = bounded
? "SELECT payload, payload_size, watermark_db_version, is_final "
"FROM cloudsync_payload_chunks WHERE since_db_version = ? AND until_db_version = ?"
: "SELECT payload, payload_size, watermark_db_version, is_final "
"FROM cloudsync_payload_chunks WHERE since_db_version = ?";
rc = sqlite3_prepare_v2(db, chunk_sql, -1, &stmt, NULL);
if (rc != SQLITE_OK) {
sqlite3_result_error(context, sqlite3_errmsg(db), -1);
sqlite3_result_error_code(context, rc);
return rc;
}
sqlite3_bind_int64(stmt, 1, db_version);
if (bounded) sqlite3_bind_int64(stmt, 2, until);
}
sqlite3_bind_int64(stmt, 1, db_version);

int64_t new_db_version = db_version;
int64_t last_optimistic_version = -1;
Expand All @@ -1936,7 +2035,7 @@ int cloudsync_network_send_changes_internal (sqlite3_context *context, int argc,
char batch_id[UUID_STR_MAXLEN];
cloudsync_uuid_v7_string(batch_id, true);

while ((rc = sqlite3_step(stmt)) == SQLITE_ROW) {
while (stmt && (rc = sqlite3_step(stmt)) == SQLITE_ROW) {
const void *blob = sqlite3_column_blob(stmt, 0);
int blob_size = sqlite3_column_bytes(stmt, 0);
int64_t payload_size = sqlite3_column_int64(stmt, 1);
Expand Down Expand Up @@ -1971,17 +2070,21 @@ int cloudsync_network_send_changes_internal (sqlite3_context *context, int argc,
sent_bytes += payload_size;
if (watermark > new_db_version) new_db_version = watermark;
}
if (rc != SQLITE_DONE) {
sqlite3_result_error(context, sqlite3_errmsg(db), -1);
sqlite3_result_error_code(context, rc);
goto cleanup;
if (stmt) {
if (rc != SQLITE_DONE) {
sqlite3_result_error(context, sqlite3_errmsg(db), -1);
sqlite3_result_error_code(context, rc);
goto cleanup;
}
sqlite3_finalize(stmt);
stmt = NULL;
}
sqlite3_finalize(stmt);
stmt = NULL;

if (!sent_any) {
// Empty local db with no server state: preserve the previous fast no-op path.
if (db_version == 0) {
// A bounded call cannot take it: an empty window says nothing about changes
// beyond the bound, and reporting "synced" there would hide the backlog.
if (db_version == 0 && !bounded) {
if (out) {
out->server_version = 0;
out->local_version = 0;
Expand Down Expand Up @@ -2018,11 +2121,18 @@ int cloudsync_network_send_changes_internal (sqlite3_context *context, int argc,
dbutils_settings_set_key_value(data, CLOUDSYNC_KEY_SEND_DBVERSION, buf);
}

// A bounded call sends a prefix of the backlog, so new_db_version stops at the bound
// and would read as "nothing left". Report the newest local version that
// network_send_window already found: network_compute_status then returns
// out-of-sync while changes remain, which is how a caller knows to call again.
int64_t local_version = new_db_version;
if (bounded && max_local > local_version) local_version = max_local;

// populate sync result
if (out) {
out->server_version = last_optimistic_version;
out->local_version = new_db_version;
out->status = network_compute_status(last_optimistic_version, last_confirmed_version, gaps_size, new_db_version);
out->local_version = local_version;
out->status = network_compute_status(last_optimistic_version, last_confirmed_version, gaps_size, local_version);
out->send_chunks = sent_chunks;
out->send_bytes = sent_bytes;
out->apply_failure_json = apply_failure_json;
Expand Down Expand Up @@ -2706,6 +2816,10 @@ int cloudsync_network_register (sqlite3 *db, char **pzErrMsg, void *ctx) {
rc = sqlite3_create_function(db, "cloudsync_network_send_changes", 0, DEFAULT_FLAGS, ctx, cloudsync_network_send_changes, NULL, NULL);
if (rc != SQLITE_OK) return rc;

// 1-argument form: cap this call's window at until_db_version.
rc = sqlite3_create_function(db, "cloudsync_network_send_changes", 1, DEFAULT_FLAGS, ctx, cloudsync_network_send_changes, NULL, NULL);
if (rc != SQLITE_OK) return rc;

rc = sqlite3_create_function(db, "cloudsync_network_receive_changes", 0, DEFAULT_FLAGS, ctx, cloudsync_network_receive_changes, NULL, NULL);
if (rc != SQLITE_OK) return rc;

Expand Down
Loading
Loading