From 09558c9fa934c46a0277daae12c6554fbfae58e0 Mon Sep 17 00:00:00 2001 From: Andrea Donetti Date: Wed, 23 Sep 2026 17:25:21 -0600 Subject: [PATCH] feat(send): bound a send to a number of local database versions cloudsync_network_send_changes() sent the whole backlog as one batch. A send is all or nothing -- the server confirms the window only once every chunk has applied, and a failed batch is re-sent whole under a new id -- so after a long offline period or a bulk import that batch can be large enough to keep failing, re-uploading everything each attempt. The new 1-argument form sends at most that many local database versions, so the upload becomes several independently confirmed batches and a failure costs one bounded window instead of the backlog. The 0-argument form is unchanged, and cloudsync_network_sync() still calls the unbounded path. The argument counts local versions rather than naming an upper bound. db_version is shared with merges, so the local backlog is sparse in that space: a bound picked without knowing which versions hold local changes can cover a range with none of them, making no progress and giving the caller no way to know how much further to go. Counting keeps the window non-empty whenever anything is pending. A bounded call also reports the newest local version as send.localVersion, so send.status stays out-of-sync while changes remain past the window. Both that value and the window bound come from one pass over the metadata tables: local rows carry site_id 0 there and db_version > since seeks the (db_version) index, so the cost follows the pending backlog. Reading them from cloudsync_changes instead would hit its site_id-only plan, a full scan that materialises every value through cloudsync_col_value(). Nothing advances the send checkpoint past a window that was never announced: later local writes are numbered below such a checkpoint and would be stranded, and the server builds its coverage from announced batches, so the skipped range would stay a permanent gap. Co-Authored-By: Claude Opus 5 (1M context) --- API.md | 20 +++++- CHANGELOG.md | 4 ++ src/network/network.c | 156 ++++++++++++++++++++++++++++++++++++------ test/network_unit.c | 121 ++++++++++++++++++++++++++++++++ 4 files changed, 279 insertions(+), 22 deletions(-) diff --git a/API.md b/API.md index 48d789c2..6322be78 100644 --- a/API.md +++ b/API.md @@ -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: @@ -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"}}}' ``` diff --git a/CHANGELOG.md b/CHANGELOG.md index e4c1938b..d7ea5aad 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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. diff --git a/src/network/network.c b/src/network/network.c index cb920f92..5df3d3a8 100644 --- a/src/network/network.c +++ b/src/network/network.c @@ -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); @@ -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; @@ -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); @@ -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; @@ -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; @@ -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; diff --git a/test/network_unit.c b/test/network_unit.c index 71d1001d..9e1ec980 100644 --- a/test/network_unit.c +++ b/test/network_unit.c @@ -662,6 +662,126 @@ static bool test_curl_async_dns(void) { } #endif +// Bounded send: cloudsync_network_send_changes(until) caps the window at until, so a +// backlog uploads as several coherent all-or-nothing batches. Records every /apply +// window so the test can assert all chunks of a batch announce the same one. +static int64_t apply_min[64], apply_max[64]; +static int napply; +static int64_t apply_optimistic; +static bool send_no_status; // status probe answers without a version, so the local fallback governs + +static NETWORK_RESULT send_responder(const char *endpoint, const char *request) { + NETWORK_RESULT r = {0}; + size_t n = endpoint ? strlen(endpoint) : 0; + bool is_apply = (n >= 6 && strcmp(endpoint + n - 6, "/apply") == 0); + if (is_apply && request && napply < 64) { + apply_min[napply] = json_int_after(request, "\"dbVersionMin\":", -1); + apply_max[napply] = json_int_after(request, "\"dbVersionMax\":", -1); + if (apply_max[napply] > apply_optimistic) apply_optimistic = apply_max[napply]; + napply++; + } + // Both /apply and the status probe answer with the same sync-status shape. + char *json = cloudsync_memory_alloc(256); + if (send_no_status && !is_apply) snprintf(json, 256, "{\"data\":{}}"); + else snprintf(json, 256, "{\"data\":{\"lastOptimisticVersion\":%lld,\"lastConfirmedVersion\":%lld,\"gaps\":[]}}", + (long long)apply_optimistic, (long long)apply_optimistic); + r.code = CLOUDSYNC_NETWORK_BUFFER; + r.buffer = json; + r.blen = strlen(json); + return r; +} + +static int64_t send_checkpoint(sqlite3 *db) { + return db_int(db, "SELECT coalesce((SELECT value FROM cloudsync_settings WHERE key='send_dbversion'),0)"); +} + +static char *send_changes(sqlite3 *db, const char *arg) { + static char out[1024]; + char sql[128]; + snprintf(sql, sizeof(sql), "SELECT cloudsync_network_send_changes(%s)", arg ? arg : ""); + sqlite3_stmt *vm = NULL; + out[0] = 0; + if (sqlite3_prepare_v2(db, sql, -1, &vm, NULL) == SQLITE_OK && sqlite3_step(vm) == SQLITE_ROW) { + const char *t = (const char *)sqlite3_column_text(vm, 0); + if (t) snprintf(out, sizeof(out), "%s", t); + } else { + snprintf(out, sizeof(out), "ERROR: %s", sqlite3_errmsg(db)); + } + sqlite3_finalize(vm); + return out; +} + +static bool test_bounded_send(void) { + bool ok = true; + sqlite3 *db = stream_db(true); + ok = ok && expect(db != NULL, "db", NULL); + if (!ok) { stream_close(db); return false; } + + // Five rows, one transaction each, so the backlog spans five local db_versions. + for (int i = 1; i <= 5 && ok; i++) { + char sql[128]; + snprintf(sql, sizeof(sql), "INSERT INTO t VALUES('r%d', 'v%d');", i, i); + ok = ok && db_exec(db, sql) == SQLITE_OK; + } + ok = ok && expect(db_int(db, "SELECT max(db_version) FROM cloudsync_changes") == 5, "five local db_versions", NULL); + + // Bad arguments are rejected before any I/O. + ok = ok && expect(strstr(send_changes(db, "'abc'"), "must be an integer") != NULL, "non-integer count rejected", NULL); + ok = ok && expect(strstr(send_changes(db, "0"), "greater than zero") != NULL, "zero count rejected", NULL); + + napply = 0; apply_optimistic = 0; + network_test_set_responder(send_responder); + + // Two db_versions per call: the window is [1,2], and every chunk of the batch + // announces it. + char *json = send_changes(db, "2"); + bool same_window = (napply > 0); + for (int i = 0; i < napply; i++) same_window = same_window && apply_min[i] == 1 && apply_max[i] == 2; + ok = ok && expect(same_window, "every chunk announces [1,2]", json); + ok = ok && expect(send_checkpoint(db) == 2, "checkpoint advanced to the bound", json); + ok = ok && expect(strstr(json, "\"status\":\"out-of-sync\"") != NULL, "backlog reported out-of-sync", json); + ok = ok && expect(strstr(json, "\"localVersion\":5") != NULL, "localVersion is the real backlog", json); + + // Asking for more versions than remain sends the rest and stops at the backlog. + napply = 0; + json = send_changes(db, "100"); + same_window = (napply > 0); + for (int i = 0; i < napply; i++) same_window = same_window && apply_min[i] == 3 && apply_max[i] == 5; + ok = ok && expect(same_window, "second batch announces [3,5]", json); + ok = ok && expect(send_checkpoint(db) == 5, "checkpoint at the end of the backlog", json); + ok = ok && expect(strstr(json, "\"status\":\"synced\"") != NULL, "drained backlog reports synced", json); + + // Sparse local versions: a merge consumes db_versions that hold no local change, so + // a count-based window has to skip over them instead of stalling on an empty range. + ok = ok && db_exec(db, "INSERT INTO t VALUES('remote1','x');" + "UPDATE \"t_cloudsync\" SET site_id=1, db_version=6 WHERE pk=cloudsync_pk_encode('remote1');") == SQLITE_OK; + ok = ok && db_exec(db, "INSERT INTO t VALUES('r6','v6');") == SQLITE_OK; + int64_t local_tail = db_int(db, "SELECT max(db_version) FROM \"t_cloudsync\" WHERE site_id=0"); + ok = ok && expect(local_tail > 6, "a local change sits past the remote-only version", NULL); + napply = 0; + json = send_changes(db, "1"); + same_window = (napply > 0); + for (int i = 0; i < napply; i++) same_window = same_window && apply_max[i] == local_tail; + ok = ok && expect(same_window, "window skips the remote-only db_version", json); + ok = ok && expect(send_checkpoint(db) == local_tail, "checkpoint reaches the local change", json); + + // Nothing pending: no batch, and the checkpoint stays put. Advancing past a range + // nobody announced would strand later local writes and leave a permanent gap in the + // server's ranges. Checked with the status probe answering without a version, so the + // local fallback governs -- otherwise the server's optimistic value masks it. + napply = 0; + send_no_status = true; + int64_t before = send_checkpoint(db); + json = send_changes(db, "10"); + ok = ok && expect(napply == 0, "nothing pending sends no batch", json); + ok = ok && expect(send_checkpoint(db) == before, "nothing pending leaves the checkpoint alone", json); + send_no_status = false; + + network_test_set_responder(NULL); + stream_close(db); + return ok; +} + int main(void) { #if !defined(_WIN32) && !defined(CLOUDSYNC_OMIT_CURL) check("HTTP deadlines and interrupt: API cap, artifact stall, cancel:", test_stalled_http_timeout()); @@ -678,6 +798,7 @@ int main(void) { check("non-buffer response is a no-op:", test_non_buffer_is_noop()); check("send batch /apply payload (window / batchId / chunkIndex / isFinal):", test_apply_json_payload_batch()); check("network_compute_status:", test_compute_status()); + check("bounded send: window, checkpoint, empty-window guard:", test_bounded_send()); check("receive stream: capped paging, failure replay, checkpoint errors:", test_stream_paging()); check("receive stream: fragmented values:", test_stream_fragments()); check("receive stream: server without watermark:", test_stream_no_watermark());