From e0bd58311fa9e20fd7c7f0edf40aa6ba7e57936a Mon Sep 17 00:00:00 2001 From: Andrea Donetti Date: Wed, 23 Sep 2026 18:03:24 -0600 Subject: [PATCH 01/14] feat(sqlite): cap the window cloudsync_payload_chunks prepares Preparation is unbounded. #75 removed the quadratic term, so a drain is now linear in the stream, but linear still means a large enough tenant cannot finish inside the job deadline -- and because nothing is persisted until the loop reaches is_final, every attempt restarts from chunk 0 and keeps no progress. max_window_bytes ends the scan once that many payload bytes have been emitted, reporting an ordinary complete stream over a smaller window: is_final with the watermark lowered to the last db_version emitted. The caller checkpoints there and asks again, so preparation is bounded with no new resumable state and no protocol change. Unset (the default) is byte-identical to today. Two conditions are load-bearing. A window may not end mid-value or mid-db_version: the receive cursor must land on a complete db_version or the next request skips the unapplied remainder, since it resumes with db_version > since and no seq. And because that boundary test is what stops the scan, a db_version larger than the whole budget is still emitted in full, so a window can never come out empty and stall the drain. window_capped is an output column, so it does not move the hidden columns a table-valued call binds positionally; max_window_bytes is declared last, taking argument 8 and leaving 1..7 as they were. Also: an explicit NULL for resume_db_version now means "not given", as it already did for site_id. Callers must pass NULLs to reach a later argument, and reading one as a resume point of 0 silently restarted the scan at the start of the window. Co-Authored-By: Claude Opus 5 (1M context) --- src/sqlite/cloudsync_sqlite.c | 72 +++++++++++++++++++++++++--- test/review_regressions.c | 88 +++++++++++++++++++++++++++++++++++ 2 files changed, 153 insertions(+), 7 deletions(-) diff --git a/src/sqlite/cloudsync_sqlite.c b/src/sqlite/cloudsync_sqlite.c index 8345dd8..4e014e3 100644 --- a/src/sqlite/cloudsync_sqlite.c +++ b/src/sqlite/cloudsync_sqlite.c @@ -1035,6 +1035,12 @@ typedef struct { int64_t next_seq; int64_t next_frag_offset; bool is_final; + // Window cap (max_window_bytes): payload bytes emitted so far by this scan, and + // whether the scan ended because the budget ran out rather than because the + // window was drained. window_bytes counts across chunks, not within one. + int64_t max_window_bytes; + int64_t window_bytes; + bool window_capped; } cloudsync_payload_chunks_cursor; static int payload_chunks_connect(sqlite3 *db, void *aux, int argc, const char *const *argv, sqlite3_vtab **vtab, char **err) { @@ -1048,7 +1054,12 @@ static int payload_chunks_connect(sqlite3 *db, void *aux, int argc, const char * // back as the resume_* inputs (cols 15..17) to continue the drain without // a spool table — O(1) seek per chunk instead of replaying from since. "next_db_version INTEGER, next_seq INTEGER, next_frag_offset INTEGER, is_final INTEGER, " - "resume_db_version HIDDEN, resume_seq HIDDEN, resume_frag_offset HIDDEN)"); + // window_capped (col 15) is an output, so it does not shift the hidden columns + // a table-valued call binds positionally. max_window_bytes is declared last for + // the same reason: it becomes argument 8, leaving arguments 1..7 as they were. + "window_capped INTEGER, " + "resume_db_version HIDDEN, resume_seq HIDDEN, resume_frag_offset HIDDEN, " + "max_window_bytes HIDDEN)"); if (rc != SQLITE_OK) return rc; cloudsync_payload_chunks_vtab *p = sqlite3_malloc64(sizeof(*p)); if (!p) return SQLITE_NOMEM; @@ -1086,9 +1097,9 @@ static int payload_chunks_best_index(sqlite3_vtab *vtab, sqlite3_index_info *idx // in a fixed order regardless of how SQLite presents constraints. idxNum bit k // is set when handled_cols[k] is bound; xFilter reads argv in this same order. // bit0=since_db_version(7) bit1=site_id(8) bit2=until_db_version(9) - // bit3=exclude_filter_site_id(10) bit4=resume_db_version(15) - // bit5=resume_seq(16) bit6=resume_frag_offset(17) - static const int handled_cols[] = {7, 8, 9, 10, 15, 16, 17}; + // bit3=exclude_filter_site_id(10) bit4=resume_db_version(16) + // bit5=resume_seq(17) bit6=resume_frag_offset(18) bit7=max_window_bytes(19) + static const int handled_cols[] = {7, 8, 9, 10, 16, 17, 18, 19}; int argv_index = 1; int idxnum = 0; for (size_t k = 0; k < sizeof(handled_cols) / sizeof(handled_cols[0]); ++k) { @@ -1305,9 +1316,38 @@ static void payload_chunks_set_next_cursor(cloudsync_payload_chunks_cursor *c) { } } +// Stop the scan once max_window_bytes is spent, at the first db_version boundary, and +// report the window as an ordinary complete stream ending at that reduced watermark. +// The caller then checkpoints there and asks again, so preparation is bounded without +// any new resumable state. +// +// Two conditions are load-bearing. Never cap mid-value (frag_active) or mid-db_version +// (next_dbv == dbv_max): the receive cursor must land on a complete db_version or the +// next request skips the unapplied remainder, since it resumes with db_version > since +// and no seq. And because the boundary test is what stops us, a db_version larger than +// the whole budget is still emitted in full -- otherwise a window could end up empty +// and the drain would never advance at all. +static void payload_chunks_apply_window_cap(cloudsync_payload_chunks_cursor *c) { + if (c->max_window_bytes <= 0 || c->eof) return; + c->window_bytes += c->payload_size; + if (c->is_final || c->frag_active) return; + if (c->window_bytes < c->max_window_bytes) return; + if (c->next_dbv == c->dbv_max) return; // still inside a db_version + + c->window_capped = true; + c->is_final = true; + c->watermark = c->dbv_max; // the window the caller checkpoints at + c->next_dbv = c->watermark; + c->next_seq = 0; + c->next_frag_offset = 0; +} + static int payload_chunks_advance(cloudsync_payload_chunks_cursor *c) { int rc = payload_chunks_build_next(c); - if (rc == SQLITE_OK && !c->eof) payload_chunks_set_next_cursor(c); + if (rc == SQLITE_OK && !c->eof) { + payload_chunks_set_next_cursor(c); + payload_chunks_apply_window_cap(c); + } return rc; } @@ -1353,9 +1393,22 @@ static int payload_chunks_filter(sqlite3_vtab_cursor *cursor, int idxnum, const } if (idxnum & 4) until = sqlite3_value_int64(argv[argi++]); if (idxnum & 8) exclude = (sqlite3_value_int(argv[argi++]) != 0); - if (idxnum & 16) { resume_dbv = sqlite3_value_int64(argv[argi++]); positional = true; } + // An explicit NULL means "not given", as it already does for site_id. Callers must + // pass NULLs to reach a later argument, and treating one as a resume point of 0 + // would silently restart the scan at the beginning of the window. + if (idxnum & 16) { + if (sqlite3_value_type(argv[argi]) != SQLITE_NULL) { + resume_dbv = sqlite3_value_int64(argv[argi]); + positional = true; + } + argi++; + } if (idxnum & 32) resume_seq = sqlite3_value_int64(argv[argi++]); if (idxnum & 64) resume_frag = sqlite3_value_int64(argv[argi++]); + if (idxnum & 128) { + int64_t cap = sqlite3_value_int64(argv[argi++]); + c->max_window_bytes = (cap > 0) ? cap : 0; // <= 0 means no cap + } // Resolve the site filter: // exclude=true -> all sites except filter_site_id (CHECK path); site required @@ -1445,7 +1498,11 @@ static int payload_chunks_filter(sqlite3_vtab_cursor *cursor, int idxnum, const } static int payload_chunks_next(sqlite3_vtab_cursor *cursor) { - return payload_chunks_advance((cloudsync_payload_chunks_cursor *)cursor); + cloudsync_payload_chunks_cursor *c = (cloudsync_payload_chunks_cursor *)cursor; + // A capped window ends the scan: the rows past it belong to the next window, and + // emitting them here would contradict the reduced watermark already reported. + if (c->window_capped) { c->eof = true; return SQLITE_OK; } + return payload_chunks_advance(c); } static int payload_chunks_eof(sqlite3_vtab_cursor *cursor) { @@ -1466,6 +1523,7 @@ static int payload_chunks_column(sqlite3_vtab_cursor *cursor, sqlite3_context *c case 12: sqlite3_result_int64(ctx, c->next_seq); break; case 13: sqlite3_result_int64(ctx, c->next_frag_offset); break; case 14: sqlite3_result_int(ctx, c->is_final ? 1 : 0); break; + case 15: sqlite3_result_int(ctx, c->window_capped ? 1 : 0); break; default: sqlite3_result_null(ctx); break; } return SQLITE_OK; diff --git a/test/review_regressions.c b/test/review_regressions.c index 14e27fd..c0ffcf2 100644 --- a/test/review_regressions.c +++ b/test/review_regressions.c @@ -690,6 +690,93 @@ static void test_block_group_atomicity(void) { } } +// A window cap must produce ordinary complete streams over smaller windows: repeated +// capped calls tile the stream exactly, every window ends on a db_version boundary, and +// the drain always advances -- a window that emitted nothing would never progress. +static void test_payload_window_cap(void) { + sqlite3 *db = open_db(); + // The smallest chunk size the setting allows, so 60 rows of 20 KB span several + // chunks and a budget has something to split. + CHECK(sql(db, "CREATE TABLE t(id TEXT PRIMARY KEY NOT NULL, v BLOB);" + "SELECT cloudsync_init('t');" + "SELECT cloudsync_set('payload_max_chunk_size','262144');") == SQLITE_OK); + // Two shapes, both needed. Single-row transactions give chunk boundaries that + // coincide with db_version boundaries, which is where a window may end and is the + // common production shape. The two 15-row transactions are ~300 KB each, larger + // than a chunk, so a chunk boundary also falls *inside* a db_version -- that is + // what exercises the rule that a window must not end there. Capping mid-version + // would leave the remainder unsent, because the next window starts past it. + int row = 0; + for (int i = 0; i < 40; i++) { + char buf[128]; + snprintf(buf, sizeof(buf), "INSERT INTO t VALUES('r%03d', randomblob(20000));", row++); + CHECK(sql(db, buf) == SQLITE_OK); + } + for (int txn = 0; txn < 2; txn++) { + CHECK(sql(db, "BEGIN;") == SQLITE_OK); + for (int i = 0; i < 15; i++) { + char buf[128]; + snprintf(buf, sizeof(buf), "INSERT INTO t VALUES('r%03d', randomblob(20000));", row++); + CHECK(sql(db, buf) == SQLITE_OK); + } + CHECK(sql(db, "COMMIT;") == SQLITE_OK); + } + + // Baseline: the whole window in one uncapped scan. + sqlite3_stmt *vm = NULL; + int64_t total_bytes = 0, total_chunks = 0, watermark = 0; + CHECK(sqlite3_prepare_v2(db, "SELECT count(*), sum(payload_size), max(watermark_db_version), " + "max(window_capped) FROM cloudsync_payload_chunks(0,NULL,NULL,false)", + -1, &vm, NULL) == SQLITE_OK); + if (sqlite3_step(vm) == SQLITE_ROW) { + total_chunks = sqlite3_column_int64(vm, 0); + total_bytes = sqlite3_column_int64(vm, 1); + watermark = sqlite3_column_int64(vm, 2); + CHECK(sqlite3_column_int(vm, 3) == 0); // nothing is capped without a budget + } + sqlite3_finalize(vm); + CHECK(total_chunks > 1 && total_bytes > 0 && watermark == 42); + + // Two budgets: one that spans several chunks, and one below a single chunk so the + // "always emit one whole db_version" guarantee is what keeps the drain moving. + const int64_t budgets[] = {200000, 1}; + for (size_t b = 0; b < sizeof(budgets) / sizeof(budgets[0]); ++b) { + int64_t since = 0, sum_bytes = 0, sum_chunks = 0; + int windows = 0; + bool capped = true; + while (capped && windows < 100) { + char q[512]; + snprintf(q, sizeof(q), + "SELECT count(*), coalesce(sum(payload_size),0), max(watermark_db_version), " + "max(window_capped), min(db_version_min) " + "FROM cloudsync_payload_chunks(%lld,NULL,NULL,false,NULL,NULL,NULL,%lld)", + (long long)since, (long long)budgets[b]); + vm = NULL; + CHECK(sqlite3_prepare_v2(db, q, -1, &vm, NULL) == SQLITE_OK); + if (sqlite3_step(vm) != SQLITE_ROW || sqlite3_column_int64(vm, 0) == 0) { sqlite3_finalize(vm); break; } + int64_t chunks = sqlite3_column_int64(vm, 0); + int64_t bytes = sqlite3_column_int64(vm, 1); + int64_t wm = sqlite3_column_int64(vm, 2); + capped = sqlite3_column_int(vm, 3) != 0; + int64_t first = sqlite3_column_int64(vm, 4); + sqlite3_finalize(vm); + + CHECK(first == since + 1); // windows abut, no gap and no overlap + CHECK(wm > since); // always advances, so the drain terminates + sum_chunks += chunks; + sum_bytes += bytes; + since = wm; + windows++; + } + CHECK(windows > 1); // the budget really did split the stream + CHECK(since == watermark); // and the last window reached the end + CHECK(sum_bytes == total_bytes); // tiles the uncapped stream exactly + CHECK(sum_chunks == total_chunks); + } + + CHECK(close_db(db) == SQLITE_OK); +} + int main(void) { CHECK(sqlite3_config(SQLITE_CONFIG_GETMALLOC, &memory) == SQLITE_OK); sqlite3_mem_methods faults = memory; @@ -712,6 +799,7 @@ int main(void) { test_block_migration_orphan(); test_block_not_null_payload(); test_block_group_atomicity(); + test_payload_window_cap(); test_refill_error(); test_block_oom(); cloudsync_memory_finalize(); From 5c55da0668e0188796e2c9a61412f166716f4d1d Mon Sep 17 00:00:00 2001 From: Andrea Donetti Date: Thu, 24 Sep 2026 11:19:58 -0600 Subject: [PATCH 02/14] feat(postgres): cap the window cloudsync_payload_chunks prepares Mirrors the SQLite side: max_window_bytes ends the scan once that many payload bytes have been emitted, at the first db_version boundary, reporting an ordinary complete stream over a smaller window -- is_final with the watermark lowered to the last db_version emitted, and window_capped to say why. Unset behaves exactly as before. The same two conditions are load-bearing here. A window may not end mid-value or mid-db_version, because the receive cursor resumes with db_version > since and no seq, so the unapplied remainder would be skipped. And because that boundary test is what stops the scan, a db_version larger than the whole budget is still emitted in full, so a window can never come out empty and stall the drain. The SQL surface changes, so EXTVERSION moves 1.1 -> 1.2 and the binary semver becomes 1.2.0 -- new optional functionality, backward compatible, which is a MINOR bump. cloudsync--1.1--1.2.sql drops the old function before creating the new one: CREATE OR REPLACE cannot change a return type, and because the new parameter has a default, keeping both would leave every existing 7-argument call matching two candidates. max_window_bytes is declared last so arguments 1..7 keep their meaning for the positional calls the /check job makes. Co-Authored-By: Claude Opus 5 (1M context) --- src/cloudsync.h | 2 +- src/postgresql/cloudsync.sql.in | 6 +- src/postgresql/cloudsync_postgresql.c | 39 +++++- .../migrations/cloudsync--1.1--1.2.sql | 46 +++++++ test/postgresql/65_payload_window_cap.sql | 121 ++++++++++++++++++ test/postgresql/full_test.sql | 1 + 6 files changed, 209 insertions(+), 6 deletions(-) create mode 100644 src/postgresql/migrations/cloudsync--1.1--1.2.sql create mode 100644 test/postgresql/65_payload_window_cap.sql diff --git a/src/cloudsync.h b/src/cloudsync.h index ffafc4b..a147605 100644 --- a/src/cloudsync.h +++ b/src/cloudsync.h @@ -18,7 +18,7 @@ extern "C" { #endif -#define CLOUDSYNC_VERSION "1.1.4" +#define CLOUDSYNC_VERSION "1.2.0" // LZ4's block format cannot expand input by more than 255:1, so a compressed payload // declaring a larger expansion is forged or corrupt (see cloudsync_payload_apply). #define CLOUDSYNC_PAYLOAD_LZ4_MAX_RATIO 255 diff --git a/src/postgresql/cloudsync.sql.in b/src/postgresql/cloudsync.sql.in index 92bb6f4..0cc4594 100644 --- a/src/postgresql/cloudsync.sql.in +++ b/src/postgresql/cloudsync.sql.in @@ -156,7 +156,8 @@ CREATE OR REPLACE FUNCTION cloudsync_payload_chunks( exclude_filter_site_id boolean DEFAULT false, resume_db_version bigint DEFAULT NULL, resume_seq bigint DEFAULT NULL, - resume_frag_offset bigint DEFAULT NULL + resume_frag_offset bigint DEFAULT NULL, + max_window_bytes bigint DEFAULT NULL ) RETURNS TABLE ( payload bytea, @@ -169,7 +170,8 @@ RETURNS TABLE ( next_db_version bigint, next_seq bigint, next_frag_offset bigint, - is_final boolean + is_final boolean, + window_capped boolean ) AS 'MODULE_PATHNAME', 'cloudsync_payload_chunks' LANGUAGE C VOLATILE; diff --git a/src/postgresql/cloudsync_postgresql.c b/src/postgresql/cloudsync_postgresql.c index 9f7145a..b05ef18 100644 --- a/src/postgresql/cloudsync_postgresql.c +++ b/src/postgresql/cloudsync_postgresql.c @@ -1083,6 +1083,12 @@ typedef struct { bool eof; int64 chunk_index; int64 watermark; + // Window cap (max_window_bytes): payload bytes emitted so far by this scan, and + // whether it ended because the budget ran out rather than because the window was + // drained. window_bytes counts across chunks, not within one. + int64 max_window_bytes; + int64 window_bytes; + bool window_capped; int max_size; int frag_target; @@ -1380,6 +1386,9 @@ Datum cloudsync_payload_chunks(PG_FUNCTION_ARGS) { int64 resume_dbv = PG_ARGISNULL(4) ? 0 : PG_GETARG_INT64(4); int64 resume_seq = PG_ARGISNULL(5) ? 0 : PG_GETARG_INT64(5); int64 resume_frag = PG_ARGISNULL(6) ? 0 : PG_GETARG_INT64(6); + // Cap the whole prepared window, not one chunk: <= 0 and NULL both mean no cap. + int64 window_cap = PG_ARGISNULL(7) ? 0 : PG_GETARG_INT64(7); + st->max_window_bytes = (window_cap > 0) ? window_cap : 0; // Site filter resolution: // exclude=true -> all sites except filter_site_id (CHECK path); site required // filter given -> only that site @@ -1472,7 +1481,9 @@ Datum cloudsync_payload_chunks(PG_FUNCTION_ARGS) { PayloadChunksState *st = (PayloadChunksState *)funcctx->user_fctx; int64 rows = 0, dbv_min = 0, dbv_max = 0; - bytea *payload = payload_chunks_build_pg_next(st, data, &rows, &dbv_min, &dbv_max); + // A capped window ends the scan: the rows past it belong to the next window, and + // emitting them here would contradict the reduced watermark already reported. + bytea *payload = st->window_capped ? NULL : payload_chunks_build_pg_next(st, data, &rows, &dbv_min, &dbv_max); if (!payload) { if (st->portal) SPI_cursor_close(st->portal); st->portal = NULL; @@ -1497,8 +1508,29 @@ Datum cloudsync_payload_chunks(PG_FUNCTION_ARGS) { next_dbv = st->watermark; next_seq = 0; next_frag = 0; is_final = true; } - Datum outvals[11]; - bool outnulls[11] = {false,false,false,false,false,false,false,false,false,false,false}; + // Stop once max_window_bytes is spent, at the first db_version boundary, and report + // the window as an ordinary complete stream ending at that reduced watermark. The + // caller checkpoints there and asks again, so preparation is bounded with no new + // resumable state. + // + // Two conditions are load-bearing. Never cap mid-value (frag_active) or + // mid-db_version (next_dbv == dbv_max): the receive cursor must land on a complete + // db_version or the next request skips the unapplied remainder, since it resumes + // with db_version > since and no seq. And because the boundary test is what stops + // us, a db_version larger than the whole budget is still emitted in full -- + // otherwise a window could come out empty and the drain would never advance. + if (st->max_window_bytes > 0) { + st->window_bytes += VARSIZE_ANY_EXHDR(payload); + if (!is_final && !st->frag_active && st->window_bytes >= st->max_window_bytes && next_dbv != dbv_max) { + st->window_capped = true; + is_final = true; + st->watermark = dbv_max; + next_dbv = st->watermark; next_seq = 0; next_frag = 0; + } + } + + Datum outvals[12]; + bool outnulls[12] = {false,false,false,false,false,false,false,false,false,false,false,false}; outvals[0] = PointerGetDatum(payload); outvals[1] = Int64GetDatum(st->chunk_index++); outvals[2] = Int64GetDatum(VARSIZE_ANY_EXHDR(payload)); @@ -1510,6 +1542,7 @@ Datum cloudsync_payload_chunks(PG_FUNCTION_ARGS) { outvals[8] = Int64GetDatum(next_seq); outvals[9] = Int64GetDatum(next_frag); outvals[10] = BoolGetDatum(is_final); + outvals[11] = BoolGetDatum(st->window_capped); HeapTuple outtup = heap_form_tuple(st->outdesc, outvals, outnulls); SRF_RETURN_NEXT(funcctx, HeapTupleGetDatum(outtup)); } diff --git a/src/postgresql/migrations/cloudsync--1.1--1.2.sql b/src/postgresql/migrations/cloudsync--1.1--1.2.sql new file mode 100644 index 0000000..45a3f00 --- /dev/null +++ b/src/postgresql/migrations/cloudsync--1.1--1.2.sql @@ -0,0 +1,46 @@ +-- CloudSync PostgreSQL extension upgrade: 1.1 -> 1.2 +-- +-- Bounds what one call to cloudsync_payload_chunks() prepares: +-- * new max_window_bytes input, declared last so the existing positional +-- arguments 1..7 keep their meaning +-- * new window_capped output, true when the scan stopped on the budget +-- rather than because the window was drained +-- +-- Both are optional: without max_window_bytes the function behaves exactly as +-- it did in 1.1. +-- +-- Run automatically by: ALTER EXTENSION cloudsync UPDATE; + +-- The old function has to go before the new one is created. CREATE OR REPLACE +-- cannot change a return type ("cannot change return type of existing +-- function"), and because the new parameter has a default, keeping both would +-- leave a 7-argument call matching two candidates -- an ambiguous function +-- call error at every existing call site. +DROP FUNCTION IF EXISTS cloudsync_payload_chunks(bigint, bytea, bigint, boolean, bigint, bigint, bigint); + +CREATE OR REPLACE FUNCTION cloudsync_payload_chunks( + since_db_version bigint DEFAULT NULL, + filter_site_id bytea DEFAULT NULL, + until_db_version bigint DEFAULT NULL, + exclude_filter_site_id boolean DEFAULT false, + resume_db_version bigint DEFAULT NULL, + resume_seq bigint DEFAULT NULL, + resume_frag_offset bigint DEFAULT NULL, + max_window_bytes bigint DEFAULT NULL +) +RETURNS TABLE ( + payload bytea, + chunk_index bigint, + payload_size bigint, + rows bigint, + db_version_min bigint, + db_version_max bigint, + watermark_db_version bigint, + next_db_version bigint, + next_seq bigint, + next_frag_offset bigint, + is_final boolean, + window_capped boolean +) +AS 'MODULE_PATHNAME', 'cloudsync_payload_chunks' +LANGUAGE C VOLATILE; diff --git a/test/postgresql/65_payload_window_cap.sql b/test/postgresql/65_payload_window_cap.sql new file mode 100644 index 0000000..858b56c --- /dev/null +++ b/test/postgresql/65_payload_window_cap.sql @@ -0,0 +1,121 @@ +-- max_window_bytes bounds what one cloudsync_payload_chunks() call prepares. A capped +-- window is an ordinary complete stream over a smaller window, so repeated calls must +-- tile the stream exactly: same chunks, same bytes, no gap and no overlap. +-- +-- Two shapes are needed in the data. Single-row transactions give chunk boundaries that +-- coincide with db_version boundaries, which is where a window may end. The large +-- transactions are bigger than one chunk, so a chunk boundary also falls inside a +-- db_version -- that is what exercises the rule that a window must not end there, +-- because the next window resumes with db_version > since and would skip the remainder. + +\set testid '65-window-cap' +\ir helper_test_init.sql + +\connect postgres +\ir helper_psql_conn_setup.sql +DROP DATABASE IF EXISTS cloudsync_test_65; +CREATE DATABASE cloudsync_test_65; + +\connect cloudsync_test_65 +\ir helper_psql_conn_setup.sql +CREATE EXTENSION IF NOT EXISTS cloudsync; +CREATE TABLE items (id TEXT PRIMARY KEY NOT NULL, v BYTEA); +SELECT cloudsync_init('items', 'CLS', 1) AS _init \gset +-- The smallest chunk size the setting allows, so the rows below span several chunks. +SELECT cloudsync_set('payload_max_chunk_size', '262144') AS _chunk \gset + +-- 40 single-row transactions, then two of 15 rows (~300 KB each, larger than a chunk). +-- \gexec runs each generated INSERT as its own command, so each is its own transaction +-- and therefore its own db_version -- a DO block would make all of them one. The values +-- are random because the payload is LZ4-compressed: a repeating pattern would collapse +-- to a few bytes and never fill a chunk. +SELECT format($f$INSERT INTO items (id, v) SELECT 'r%s', (SELECT decode(string_agg(md5(random()::text), ''), 'hex') FROM generate_series(1, 1250))$f$, i) + FROM generate_series(1, 40) i \gexec + +BEGIN; +INSERT INTO items (id, v) + SELECT 'b1_' || i, (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 1250) g) + FROM generate_series(1, 15) i; +COMMIT; +BEGIN; +INSERT INTO items (id, v) + SELECT 'b2_' || i, (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 1250) g) + FROM generate_series(1, 15) i; +COMMIT; + +-- Baseline: the whole window, uncapped. Nothing is capped without a budget. +SELECT count(*) AS base_chunks, sum(payload_size) AS base_bytes, + max(watermark_db_version) AS base_wm, + (count(*) > 1 AND sum(payload_size) > 0 AND bool_or(window_capped) IS FALSE) AS base_ok + FROM cloudsync_payload_chunks(0, NULL, NULL, false) \gset +\if :base_ok +\echo [PASS] (:testid) uncapped baseline spans several chunks and reports no cap +\else +\echo [FAIL] (:testid) baseline chunks=:base_chunks bytes=:base_bytes +SELECT (:fail::int + 1) AS fail \gset +\endif + +-- Drain the same stream in capped windows, following watermark_db_version each time, +-- and add up what the windows covered. A budget below one chunk also proves the drain +-- still advances: a db_version larger than the whole budget must be emitted in full. +CREATE FUNCTION pg_temp.drain_capped(cap bigint) +RETURNS TABLE (windows int, chunks bigint, bytes bigint, last_wm bigint, contiguous boolean) AS $$ +DECLARE + since bigint := 0; + w int := 0; + c bigint := 0; + b bigint := 0; + ok boolean := true; + r record; +BEGIN + LOOP + SELECT count(*) AS n, coalesce(sum(payload_size), 0) AS sz, + max(watermark_db_version) AS wm, bool_or(window_capped) AS capped, + min(db_version_min) AS first_dbv + INTO r + FROM cloudsync_payload_chunks(since, NULL, NULL, false, NULL, NULL, NULL, cap); + EXIT WHEN r.n = 0; + -- Windows must abut exactly, and each must advance or the drain never ends. + IF r.first_dbv <> since + 1 OR r.wm <= since THEN ok := false; END IF; + w := w + 1; c := c + r.n; b := b + r.sz; since := r.wm; + EXIT WHEN NOT r.capped OR w > 100; + END LOOP; + RETURN QUERY SELECT w, c, b, since, ok; +END $$ LANGUAGE plpgsql; + +SELECT d.windows AS w1, d.chunks AS c1, d.bytes AS b1, d.last_wm AS wm1, + (d.windows > 1 AND d.contiguous AND d.chunks = :base_chunks::bigint + AND d.bytes = :base_bytes::bigint AND d.last_wm = :base_wm::bigint) AS cap1_ok + FROM pg_temp.drain_capped(200000) d \gset +\if :cap1_ok +\echo [PASS] (:testid) a 200 KB budget splits the stream into :w1 windows that tile it exactly +\else +\echo [FAIL] (:testid) windows=:w1 chunks=:c1/:base_chunks bytes=:b1/:base_bytes wm=:wm1/:base_wm +SELECT (:fail::int + 1) AS fail \gset +\endif + +SELECT d.windows AS w2, d.chunks AS c2, d.bytes AS b2, d.last_wm AS wm2, + (d.windows > 1 AND d.contiguous AND d.chunks = :base_chunks::bigint + AND d.bytes = :base_bytes::bigint AND d.last_wm = :base_wm::bigint) AS cap2_ok + FROM pg_temp.drain_capped(1) d \gset +\if :cap2_ok +\echo [PASS] (:testid) a 1-byte budget still advances and tiles the stream exactly +\else +\echo [FAIL] (:testid) windows=:w2 chunks=:c2/:base_chunks bytes=:b2/:base_bytes wm=:wm2/:base_wm +SELECT (:fail::int + 1) AS fail \gset +\endif + +-- The positional resume still works alongside the new argument. +SELECT count(*) AS res_chunks, min(db_version_min) AS res_first, + (count(*) > 0 AND min(db_version_min) = 5) AS res_ok + FROM cloudsync_payload_chunks(NULL, NULL, :base_wm, false, 5, 0, 0) \gset +\if :res_ok +\echo [PASS] (:testid) positional resume is unaffected by the new argument +\else +\echo [FAIL] (:testid) positional resume returned :res_chunks chunk(s) starting at :res_first +SELECT (:fail::int + 1) AS fail \gset +\endif + +\connect postgres +\ir helper_psql_conn_setup.sql +DROP DATABASE IF EXISTS cloudsync_test_65; diff --git a/test/postgresql/full_test.sql b/test/postgresql/full_test.sql index fe21bd0..6baafc1 100644 --- a/test/postgresql/full_test.sql +++ b/test/postgresql/full_test.sql @@ -72,6 +72,7 @@ \ir 62_deferred_fk_caller_commit.sql \ir 63_deep_savepoints.sql \ir 64_block_rewrite_leftovers.sql +\ir 65_payload_window_cap.sql \ir 66_db_version_per_transaction.sql -- 'Test summary' From b42d8e6c2994b6b320c7f128571704439780c673 Mon Sep 17 00:00:00 2001 From: Andrea Donetti Date: Thu, 24 Sep 2026 11:34:49 -0600 Subject: [PATCH 03/14] docs: document the chunk window cap and the resume arguments Adds max_window_bytes and window_capped to API.md, and the resume_* trio alongside them: those have been callable since the chunked path shipped in #50 but were never written down, which left a stateless caller no documented way to page a stream. CHANGELOG records two things beyond the feature. An explicit NULL for resume_db_version now means "not given" on SQLite, as it always has on PostgreSQL -- it used to read as a resume point of 0 and silently ignore since_db_version, which is easy to hit now that reaching a later argument means passing NULLs for the earlier ones. And the PostgreSQL extension version moves to 1.2, so deployments need ALTER EXTENSION cloudsync UPDATE. Co-Authored-By: Claude Opus 5 (1M context) --- API.md | 23 +++++++++++++++-------- CHANGELOG.md | 6 ++++++ 2 files changed, 21 insertions(+), 8 deletions(-) diff --git a/API.md b/API.md index 6322be7..30a998d 100644 --- a/API.md +++ b/API.md @@ -30,7 +30,7 @@ This document provides a reference for the SQL functions provided by the `sqlite - [Payload Functions](#payload-functions) - [`cloudsync_payload_encode()`](#cloudsync_payload_encodetbl-pk-col_name-col_value-col_version-db_version-site_id-cl-seq) - [`cloudsync_payload_blob_checked()`](#cloudsync_payload_blob_checkedsince_db_version-since_seq-filter_site_id-exclude_filter_site_id-max_estimated_payload_size) - - [`cloudsync_payload_chunks()`](#cloudsync_payload_chunkssince_db_version-filter_site_id-until_db_version-exclude_filter_site_id) + - [`cloudsync_payload_chunks()`](#cloudsync_payload_chunkssince_db_version-filter_site_id-until_db_version-exclude_filter_site_id-resume_db_version-resume_seq-resume_frag_offset-max_window_bytes) - [`cloudsync_payload_apply()`](#cloudsync_payload_applypayload) - [Network Functions](#network-functions) - [`cloudsync_network_init()`](#cloudsync_network_initmanageddatabaseid) @@ -57,7 +57,7 @@ The following settings are supported: | Key | Description | Default | Minimum | Maximum | |---|---|---:|---:|---:| -| `payload_max_chunk_size` | Maximum transport payload size generated by [`cloudsync_payload_chunks()`](#cloudsync_payload_chunkssince_db_version-filter_site_id-until_db_version-exclude_filter_site_id). Values outside the range are clamped. | `5242880` (5 MB) | `262144` (256 KB) | `33554432` (32 MB) | +| `payload_max_chunk_size` | Maximum transport payload size generated by [`cloudsync_payload_chunks()`](#cloudsync_payload_chunkssince_db_version-filter_site_id-until_db_version-exclude_filter_site_id-resume_db_version-resume_seq-resume_frag_offset-max_window_bytes). Values outside the range are clamped. | `5242880` (5 MB) | `262144` (256 KB) | `33554432` (32 MB) | | `network_connect_timeout` | Seconds allowed to connect, name resolution included. Never longer than the request's own deadline. | `30` | `1` | `86400` | | `network_request_timeout` | Seconds allowed for a whole API request (check, upload URL, apply, status). | `300` | `1` | `86400` | | `network_artifact_timeout` | Absolute limit, in seconds, for a payload upload or download. | `3600` | `1` | `86400` | @@ -436,7 +436,7 @@ SELECT cloudsync_uuid_text(cloudsync_siteid(), false); -- 0190a1b2c3d47e5f8a9b ### `cloudsync_uuid_blob(uuid)` -**Description:** Converts a UUID string into its 16-byte binary form. This is the inverse of [`cloudsync_uuid_text()`](#cloudsync_uuid_textuuid-dash_format) and lets string-based callers (for example, an HTTP `/check` endpoint holding a stringified `site_id`) pass a `site_id` to [`cloudsync_payload_chunks()`](#cloudsync_payload_chunkssince_db_version-filter_site_id-until_db_version-exclude_filter_site_id). +**Description:** Converts a UUID string into its 16-byte binary form. This is the inverse of [`cloudsync_uuid_text()`](#cloudsync_uuid_textuuid-dash_format) and lets string-based callers (for example, an HTTP `/check` endpoint holding a stringified `site_id`) pass a `site_id` to [`cloudsync_payload_chunks()`](#cloudsync_payload_chunkssince_db_version-filter_site_id-until_db_version-exclude_filter_site_id-resume_db_version-resume_seq-resume_frag_offset-max_window_bytes). **Parameters:** @@ -505,7 +505,7 @@ SELECT cloudsync_commit_alter('my_table'); **Description:** Encodes rows from `cloudsync_changes` into a single monolithic payload. This is the legacy payload API and remains fully supported for backward compatibility. -Use this API when the expected payload size is modest or when you need to interoperate with callers that expect a single BLOB. For large rowsets or large individual BLOB/TEXT values, prefer [`cloudsync_payload_chunks()`](#cloudsync_payload_chunkssince_db_version-filter_site_id-until_db_version-exclude_filter_site_id), which splits transport payloads according to `payload_max_chunk_size`. +Use this API when the expected payload size is modest or when you need to interoperate with callers that expect a single BLOB. For large rowsets or large individual BLOB/TEXT values, prefer [`cloudsync_payload_chunks()`](#cloudsync_payload_chunkssince_db_version-filter_site_id-until_db_version-exclude_filter_site_id-resume_db_version-resume_seq-resume_frag_offset-max_window_bytes), which splits transport payloads according to `payload_max_chunk_size`. **Parameters:** The function is an aggregate over the columns returned by `cloudsync_changes`: @@ -567,7 +567,7 @@ SELECT cloudsync_payload_blob_checked( --- -### `cloudsync_payload_chunks([since_db_version], [filter_site_id], [until_db_version], [exclude_filter_site_id])` +### `cloudsync_payload_chunks([since_db_version], [filter_site_id], [until_db_version], [exclude_filter_site_id], [resume_db_version], [resume_seq], [resume_frag_offset], [max_window_bytes])` **Description:** Generates sync payloads as a stream of transport-sized chunks. It is the chunk-aware evolution of [`cloudsync_payload_encode()`](#cloudsync_payload_encodetbl-pk-col_name-col_value-col_version-db_version-site_id-cl-seq), designed for large rowsets and for single BLOB/TEXT values that are larger than the configured chunk size. @@ -589,6 +589,10 @@ When a single encoded column value does not fit in one chunk, CloudSync transpar - `filter_site_id` (BLOB, optional): Site ID to filter on. With `exclude_filter_site_id` unset/`false` it selects changes **from** this site; with `exclude_filter_site_id` `true` it selects changes from every site **except** this one. If omitted (and not excluding), CloudSync uses the local site ID. - `until_db_version` (INTEGER/BIGINT, optional): Upper watermark to include. If omitted or `0`, CloudSync captures the current maximum source database version before streaming chunks. - `exclude_filter_site_id` (BOOLEAN, optional, default `false`): When `true`, stream changes from all sites **except** `filter_site_id`. This is what the `/check` download path needs — a peer must not receive its own changes back. Setting it `true` without a `filter_site_id` is an error. The site_id stored in `cloudsync_changes` is the 16-byte binary UUID; string callers can convert with [`cloudsync_uuid_blob()`](#cloudsync_uuid_blobuuid). +- `resume_db_version`, `resume_seq`, `resume_frag_offset` (INTEGER/BIGINT, optional): Resume a stream where a previous call stopped, instead of replaying it from `since_db_version`. Pass back the `next_db_version`, `next_seq` and `next_frag_offset` of the last chunk received, together with the `watermark_db_version` as `until_db_version` so the window stays fixed. This lets a stateless caller page one chunk per request without holding a cursor, and the resume is a seek rather than a re-scan. Supplying `resume_db_version` replaces `since_db_version` as the lower bound; leave all three unset (or `NULL`) for a fresh stream. +- `max_window_bytes` (INTEGER/BIGINT, optional): Cap the whole stream at roughly this many payload bytes, instead of preparing the entire window. When the budget is spent the stream ends at the next complete `db_version`, reporting `is_final` with `watermark_db_version` lowered to that point and `window_capped` set. The result is an ordinary *complete* stream over a smaller window: checkpoint at the reported watermark and call again with that value as `since_db_version` to continue. Values of `0` or `NULL` mean no cap, which is the default. + + Use it to bound how long one preparation takes on a large history, and note two properties. A window always ends on a `db_version` boundary, so a transaction larger than the budget is still emitted whole and the budget is an approximate floor rather than a hard ceiling. And because a window never ends empty, repeated calls always make progress. **Returns:** A rowset with one row per chunk: @@ -600,7 +604,10 @@ When a single encoded column value does not fit in one chunk, CloudSync transpar | `rows` | Number of encoded payload rows in this chunk. Fragment chunks usually contain one fragment row. | | `db_version_min` | Minimum source `db_version` represented by this chunk. | | `db_version_max` | Maximum source `db_version` represented by this chunk. | -| `watermark_db_version` | Stable upper watermark captured for this chunk stream. Store this after all chunks are durably transferred/applied. | +| `watermark_db_version` | Stable upper watermark captured for this chunk stream. Store this after all chunks are durably transferred/applied. With `max_window_bytes`, the final chunk reports the *reduced* watermark the capped window reached. | +| `next_db_version`, `next_seq`, `next_frag_offset` | Resume point immediately after this chunk. Pass them back as the `resume_*` arguments to continue the stream. | +| `is_final` | True on the last chunk of the stream. A capped window is a complete stream, so its last chunk reports `is_final` too. | +| `window_capped` | True when the stream ended because `max_window_bytes` was reached rather than because the window was drained — i.e. more changes exist past `watermark_db_version`. | **SQLite usage:** `cloudsync_payload_chunks` is exposed as a virtual table with hidden constraint columns: @@ -660,7 +667,7 @@ On PostgreSQL, apply chunks as individual statements from the transport/client l - Legacy payloads generated by older SQLite Sync versions. - Monolithic payloads generated by [`cloudsync_payload_encode()`](#cloudsync_payload_encodetbl-pk-col_name-col_value-col_version-db_version-site_id-cl-seq). -- Chunk-fragment payloads generated by [`cloudsync_payload_chunks()`](#cloudsync_payload_chunkssince_db_version-filter_site_id-until_db_version-exclude_filter_site_id). +- Chunk-fragment payloads generated by [`cloudsync_payload_chunks()`](#cloudsync_payload_chunkssince_db_version-filter_site_id-until_db_version-exclude_filter_site_id-resume_db_version-resume_seq-resume_frag_offset-max_window_bytes). When a v3 fragment payload is received, CloudSync stores the fragment in an internal table and returns after applying zero or more completed values. Once the final fragment for a value is received, the completed value is validated and applied. Fragments can arrive in any order, and duplicate fragment delivery is idempotent. Applying a fragment never moves the receive checkpoint. On PostgreSQL, pieces of one value applied by concurrent transactions wait for each other under `READ COMMITTED`, and fail with a retryable serialization error under `SERIALIZABLE` when they conflict; a fragment is refused under `REPEATABLE READ`, where a transaction could miss a piece committed while it waited. @@ -779,7 +786,7 @@ This means: if you get JSON back, the server was reachable and the network proto **Description:** Sends all unsent local changes to the remote server. -The send path streams payloads through [`cloudsync_payload_chunks()`](#cloudsync_payload_chunkssince_db_version-filter_site_id-until_db_version-exclude_filter_site_id), so `payload_max_chunk_size` also limits the payloads generated for network transport. Each generated chunk is uploaded/applied independently; the local send checkpoint is advanced only after the chunk stream completes successfully. +The send path streams payloads through [`cloudsync_payload_chunks()`](#cloudsync_payload_chunkssince_db_version-filter_site_id-until_db_version-exclude_filter_site_id-resume_db_version-resume_seq-resume_frag_offset-max_window_bytes), so `payload_max_chunk_size` also limits the payloads generated for network transport. Each generated chunk is uploaded/applied independently; the local send checkpoint is advanced only after the chunk stream completes successfully. 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. diff --git a/CHANGELOG.md b/CHANGELOG.md index 7e1b43a..fb0d3a5 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,12 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/). ### 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. +- **`cloudsync_payload_chunks()` can bound how much one call prepares.** Preparing a large history is unbounded work: a big enough tenant cannot finish inside a caller's time budget, and because nothing is durable until the stream reports `is_final`, an attempt that runs out of time keeps no progress and the next one restarts from the first chunk. The new `max_window_bytes` argument ends the stream once roughly that many payload bytes have been emitted, at the next complete database version. What comes back is an ordinary *complete* stream over a smaller window — `is_final` with `watermark_db_version` lowered to that point — so the caller checkpoints there and calls again to continue, with no resumable state to keep anywhere. A new `window_capped` output says the stream stopped on the budget rather than because the window was drained, so more changes exist past the watermark. Two properties are worth knowing: a window always ends on a database version boundary, so a transaction larger than the budget is still emitted whole and the budget is an approximate floor rather than a hard ceiling; and a window never ends empty, so repeated calls always make progress. Unset, the function behaves exactly as before. + +### Changed + +- **SQLite: an explicit `NULL` for `resume_db_version` on `cloudsync_payload_chunks()` now means "not given"**, matching what it has always meant for `filter_site_id` and on PostgreSQL. It was previously read as a resume point of database version 0, which silently ignored `since_db_version` and restarted the scan at the beginning of the window. Reaching a later argument requires passing `NULL` for the ones before it, so this is easy to hit: `cloudsync_payload_chunks(100, NULL, NULL, false, NULL, NULL, NULL, 1048576)` used to replay from the start of the history instead of resuming after version 100. +- **The PostgreSQL extension version moves to `1.2`.** `cloudsync_payload_chunks()` gained an argument and an output column, so existing deployments need `ALTER EXTENSION cloudsync UPDATE;` after installing the new binary. The upgrade script replaces the function: a `CREATE OR REPLACE` cannot change a return type, and leaving the old seven-argument version in place would make every existing call ambiguous against the new eight-argument one. ### Fixed From d8cd65c00055fcb310cbd6d3941187f4bff1c551 Mon Sep 17 00:00:00 2001 From: Andrea Donetti Date: Thu, 24 Sep 2026 12:19:03 -0600 Subject: [PATCH 04/14] fix: end a chunk at a version boundary once the window budget is spent The cap only looked at whether a finished chunk happened to end on a db_version boundary. When a transaction's rows and a chunk's capacity stay out of step, no chunk ever does: every one ends a row or two inside a version, the check postpones, and the cap never fires at all. Reproduced with 256 KB chunks, 20 KB rows and transactions of 12 then 13 rows, where a 200 KB budget emitted the whole 13 MB history with window_capped = 0. The overshoot is not bounded by one transaction, so a large history can still blow a preparation deadline -- exactly what the cap exists to prevent. The chunk builder now stops adding rows at the first db_version boundary once the budget is spent, so the boundary the cap needs is produced rather than waited for. Both backends. Reaching that boundary can end a chunk before it is full, so a capped drain packs the same changes into slightly more chunks, each with its own header. The tests therefore assert conservation of payload *rows*, not of bytes or chunk count, plus a per-window size bound that a never-firing cap cannot satisfy. Both suites gain the out-of-step transaction shape; without it the bug reproduces in neither. The window is also now measured in the unit max_size uses -- encoded bytes before compression -- so a budget and a chunk size mean the same thing rather than one counting compressed bytes and the other not. Co-Authored-By: Claude Opus 5 (1M context) --- API.md | 4 +- CHANGELOG.md | 2 +- src/postgresql/cloudsync_postgresql.c | 11 ++++- src/sqlite/cloudsync_sqlite.c | 11 ++++- test/postgresql/65_payload_window_cap.sql | 41 +++++++++++------- test/review_regressions.c | 51 ++++++++++++++++------- 6 files changed, 84 insertions(+), 36 deletions(-) diff --git a/API.md b/API.md index 30a998d..b08c326 100644 --- a/API.md +++ b/API.md @@ -592,7 +592,7 @@ When a single encoded column value does not fit in one chunk, CloudSync transpar - `resume_db_version`, `resume_seq`, `resume_frag_offset` (INTEGER/BIGINT, optional): Resume a stream where a previous call stopped, instead of replaying it from `since_db_version`. Pass back the `next_db_version`, `next_seq` and `next_frag_offset` of the last chunk received, together with the `watermark_db_version` as `until_db_version` so the window stays fixed. This lets a stateless caller page one chunk per request without holding a cursor, and the resume is a seek rather than a re-scan. Supplying `resume_db_version` replaces `since_db_version` as the lower bound; leave all three unset (or `NULL`) for a fresh stream. - `max_window_bytes` (INTEGER/BIGINT, optional): Cap the whole stream at roughly this many payload bytes, instead of preparing the entire window. When the budget is spent the stream ends at the next complete `db_version`, reporting `is_final` with `watermark_db_version` lowered to that point and `window_capped` set. The result is an ordinary *complete* stream over a smaller window: checkpoint at the reported watermark and call again with that value as `since_db_version` to continue. Values of `0` or `NULL` mean no cap, which is the default. - Use it to bound how long one preparation takes on a large history, and note two properties. A window always ends on a `db_version` boundary, so a transaction larger than the budget is still emitted whole and the budget is an approximate floor rather than a hard ceiling. And because a window never ends empty, repeated calls always make progress. + Use it to bound how long one preparation takes on a large history, and note three properties. A window always ends on a `db_version` boundary, so a transaction larger than the budget is still emitted whole and the budget is an approximate floor rather than a hard ceiling. Reaching that boundary can mean ending a chunk before it is full, so a capped drain packs the same changes into slightly more chunks than an uncapped one. And because a window never ends empty, repeated calls always make progress. **Returns:** A rowset with one row per chunk: @@ -604,7 +604,7 @@ When a single encoded column value does not fit in one chunk, CloudSync transpar | `rows` | Number of encoded payload rows in this chunk. Fragment chunks usually contain one fragment row. | | `db_version_min` | Minimum source `db_version` represented by this chunk. | | `db_version_max` | Maximum source `db_version` represented by this chunk. | -| `watermark_db_version` | Stable upper watermark captured for this chunk stream. Store this after all chunks are durably transferred/applied. With `max_window_bytes`, the final chunk reports the *reduced* watermark the capped window reached. | +| `watermark_db_version` | Stable upper watermark captured for this chunk stream. Store this after all chunks are durably transferred/applied. With `max_window_bytes`, **only the final chunk's value is the one to store**: earlier chunks of a capped window still report the window's provisional upper bound, because the cap is not known until the stream ends. | | `next_db_version`, `next_seq`, `next_frag_offset` | Resume point immediately after this chunk. Pass them back as the `resume_*` arguments to continue the stream. | | `is_final` | True on the last chunk of the stream. A capped window is a complete stream, so its last chunk reports `is_final` too. | | `window_capped` | True when the stream ended because `max_window_bytes` was reached rather than because the window was drained — i.e. more changes exist past `watermark_db_version`. | diff --git a/CHANGELOG.md b/CHANGELOG.md index fb0d3a5..435ec96 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,7 +9,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/). ### 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. -- **`cloudsync_payload_chunks()` can bound how much one call prepares.** Preparing a large history is unbounded work: a big enough tenant cannot finish inside a caller's time budget, and because nothing is durable until the stream reports `is_final`, an attempt that runs out of time keeps no progress and the next one restarts from the first chunk. The new `max_window_bytes` argument ends the stream once roughly that many payload bytes have been emitted, at the next complete database version. What comes back is an ordinary *complete* stream over a smaller window — `is_final` with `watermark_db_version` lowered to that point — so the caller checkpoints there and calls again to continue, with no resumable state to keep anywhere. A new `window_capped` output says the stream stopped on the budget rather than because the window was drained, so more changes exist past the watermark. Two properties are worth knowing: a window always ends on a database version boundary, so a transaction larger than the budget is still emitted whole and the budget is an approximate floor rather than a hard ceiling; and a window never ends empty, so repeated calls always make progress. Unset, the function behaves exactly as before. +- **`cloudsync_payload_chunks()` can bound how much one call prepares.** Preparing a large history is unbounded work: a big enough tenant cannot finish inside a caller's time budget, and because nothing is durable until the stream reports `is_final`, an attempt that runs out of time keeps no progress and the next one restarts from the first chunk. The new `max_window_bytes` argument ends the stream once roughly that many payload bytes have been emitted, at the next complete database version. What comes back is an ordinary *complete* stream over a smaller window — `is_final` with `watermark_db_version` lowered to that point — so the caller checkpoints there and calls again to continue, with no resumable state to keep anywhere. A new `window_capped` output says the stream stopped on the budget rather than because the window was drained, so more changes exist past the watermark. Two properties are worth knowing: a window always ends on a database version boundary, so a transaction larger than the budget is still emitted whole and the budget is an approximate floor rather than a hard ceiling; and a window never ends empty, so repeated calls always make progress. Reaching that boundary can require ending a chunk before it is full, so a capped drain packs the same changes into slightly more chunks than an uncapped one. Unset, the function behaves exactly as before. ### Changed diff --git a/src/postgresql/cloudsync_postgresql.c b/src/postgresql/cloudsync_postgresql.c index b05ef18..efb0dd1 100644 --- a/src/postgresql/cloudsync_postgresql.c +++ b/src/postgresql/cloudsync_postgresql.c @@ -1323,6 +1323,13 @@ static bytea *payload_chunks_build_pg_next(PayloadChunksState *st, cloudsync_con } if (cloudsync_payload_context_nrows(payload) > 0 && cloudsync_payload_context_bused(payload) + row_size > (size_t)st->max_size) break; + // Once the budget is spent, end the chunk at the first db_version boundary so + // the window can be capped there. Waiting for a chunk to happen to end on a + // boundary is not enough: when a transaction's rows and a chunk's capacity stay + // out of step, every chunk ends mid-version and the cap never fires at all. + if (st->max_window_bytes > 0 && cloudsync_payload_context_nrows(payload) > 0 && + st->db_version != *dbv_max && + st->window_bytes + (int64)cloudsync_payload_context_bused(payload) >= st->max_window_bytes) break; pgvalue_t *vals[9] = {0}; text *owned_texts[2] = {0}; @@ -1340,6 +1347,9 @@ static bytea *payload_chunks_build_pg_next(PayloadChunksState *st, cloudsync_con cloudsync_memory_free(payload); return NULL; } + // Measure the window in the unit max_size is expressed in -- encoded bytes before + // compression -- so a budget and a chunk size mean the same thing to a caller. + st->window_bytes += (int64)cloudsync_payload_context_bused(payload); int rc = cloudsync_payload_encode_final(payload, data); if (rc != DBRES_OK) ereport(ERROR, (errcode(cloudsync_error_sqlstate(data)), errmsg("%s", cloudsync_errmsg(data)))); int64 blob_size = 0; @@ -1520,7 +1530,6 @@ Datum cloudsync_payload_chunks(PG_FUNCTION_ARGS) { // us, a db_version larger than the whole budget is still emitted in full -- // otherwise a window could come out empty and the drain would never advance. if (st->max_window_bytes > 0) { - st->window_bytes += VARSIZE_ANY_EXHDR(payload); if (!is_final && !st->frag_active && st->window_bytes >= st->max_window_bytes && next_dbv != dbv_max) { st->window_capped = true; is_final = true; diff --git a/src/sqlite/cloudsync_sqlite.c b/src/sqlite/cloudsync_sqlite.c index 4e014e3..9aea12f 100644 --- a/src/sqlite/cloudsync_sqlite.c +++ b/src/sqlite/cloudsync_sqlite.c @@ -1272,6 +1272,13 @@ static int payload_chunks_build_next(cloudsync_payload_chunks_cursor *c) { } if (cloudsync_payload_context_nrows(payload) > 0 && cloudsync_payload_context_bused(payload) + row_size > (size_t)max_size) break; + // Once the budget is spent, end the chunk at the first db_version boundary so + // the window can be capped there. Waiting for a chunk to happen to end on a + // boundary is not enough: when a transaction's rows and a chunk's capacity stay + // out of step, every chunk ends mid-version and the cap never fires at all. + if (c->max_window_bytes > 0 && cloudsync_payload_context_nrows(payload) > 0 && + sqlite3_column_int64(c->src, 5) != c->dbv_max && + c->window_bytes + (int64_t)cloudsync_payload_context_bused(payload) >= c->max_window_bytes) break; rc = cloudsync_payload_encode_step(payload, data, 9, (dbvalue_t **)rowv); if (rc != SQLITE_OK) { cloudsync_memory_free(payload); return rc; } int64_t dbv = sqlite3_column_int64(c->src, 5); @@ -1284,6 +1291,9 @@ static int payload_chunks_build_next(cloudsync_payload_chunks_cursor *c) { if (cloudsync_payload_context_nrows(payload) == 0) { cloudsync_memory_free(payload); c->eof = true; return SQLITE_OK; } rc = cloudsync_payload_encode_final(payload, data); if (rc != SQLITE_OK) { cloudsync_memory_free(payload); return rc; } + // Measure the window in the unit max_size is expressed in -- encoded bytes before + // compression -- so a budget and a chunk size mean the same thing to a caller. + c->window_bytes += (int64_t)cloudsync_payload_context_bused(payload); c->payload = cloudsync_payload_blob(payload, &c->payload_size, &c->rows); cloudsync_memory_free(payload); c->chunk_index++; @@ -1329,7 +1339,6 @@ static void payload_chunks_set_next_cursor(cloudsync_payload_chunks_cursor *c) { // and the drain would never advance at all. static void payload_chunks_apply_window_cap(cloudsync_payload_chunks_cursor *c) { if (c->max_window_bytes <= 0 || c->eof) return; - c->window_bytes += c->payload_size; if (c->is_final || c->frag_active) return; if (c->window_bytes < c->max_window_bytes) return; if (c->next_dbv == c->dbv_max) return; // still inside a db_version diff --git a/test/postgresql/65_payload_window_cap.sql b/test/postgresql/65_payload_window_cap.sql index 858b56c..7d29479 100644 --- a/test/postgresql/65_payload_window_cap.sql +++ b/test/postgresql/65_payload_window_cap.sql @@ -43,15 +43,22 @@ INSERT INTO items (id, v) FROM generate_series(1, 15) i; COMMIT; +-- The shape that matters most: ~13 rows fill a 256 KB chunk, so transactions of 12 then +-- 13 rows keep every chunk boundary one row inside a db_version. Waiting for a chunk to +-- end on a boundary never succeeds here, so a cap that only checks after a chunk is +-- built never fires and the whole history comes out in one window. +SELECT format($f$INSERT INTO items (id, v) SELECT 'd%s_' || i, (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 1250) g) FROM generate_series(1, %s) i$f$, t, CASE WHEN t = 1 THEN 12 ELSE 13 END) + FROM generate_series(1, 20) t \gexec + -- Baseline: the whole window, uncapped. Nothing is capped without a budget. -SELECT count(*) AS base_chunks, sum(payload_size) AS base_bytes, +SELECT sum(rows) AS base_rows, sum(payload_size) AS base_bytes, max(watermark_db_version) AS base_wm, - (count(*) > 1 AND sum(payload_size) > 0 AND bool_or(window_capped) IS FALSE) AS base_ok + (count(*) > 1 AND sum(rows) > 0 AND bool_or(window_capped) IS FALSE) AS base_ok FROM cloudsync_payload_chunks(0, NULL, NULL, false) \gset \if :base_ok \echo [PASS] (:testid) uncapped baseline spans several chunks and reports no cap \else -\echo [FAIL] (:testid) baseline chunks=:base_chunks bytes=:base_bytes +\echo [FAIL] (:testid) baseline rows=:base_rows bytes=:base_bytes SELECT (:fail::int + 1) AS fail \gset \endif @@ -59,7 +66,7 @@ SELECT (:fail::int + 1) AS fail \gset -- and add up what the windows covered. A budget below one chunk also proves the drain -- still advances: a db_version larger than the whole budget must be emitted in full. CREATE FUNCTION pg_temp.drain_capped(cap bigint) -RETURNS TABLE (windows int, chunks bigint, bytes bigint, last_wm bigint, contiguous boolean) AS $$ +RETURNS TABLE (windows int, nrows bigint, maxbytes bigint, last_wm bigint, contiguous boolean) AS $$ DECLARE since bigint := 0; w int := 0; @@ -69,39 +76,41 @@ DECLARE r record; BEGIN LOOP - SELECT count(*) AS n, coalesce(sum(payload_size), 0) AS sz, - max(watermark_db_version) AS wm, bool_or(window_capped) AS capped, - min(db_version_min) AS first_dbv + SELECT count(*) AS n, coalesce(sum(rows), 0) AS nr, coalesce(sum(payload_size), 0) AS sz, + max(watermark_db_version) FILTER (WHERE is_final) AS wm, + bool_or(window_capped) AS capped, min(db_version_min) AS first_dbv INTO r FROM cloudsync_payload_chunks(since, NULL, NULL, false, NULL, NULL, NULL, cap); EXIT WHEN r.n = 0; -- Windows must abut exactly, and each must advance or the drain never ends. IF r.first_dbv <> since + 1 OR r.wm <= since THEN ok := false; END IF; - w := w + 1; c := c + r.n; b := b + r.sz; since := r.wm; + -- A window ends at the first db_version boundary at or after the budget, so it can + -- overshoot by at most the chunk that crossed it plus the version in progress. + w := w + 1; c := c + r.nr; b := greatest(b, r.sz); since := r.wm; EXIT WHEN NOT r.capped OR w > 100; END LOOP; RETURN QUERY SELECT w, c, b, since, ok; END $$ LANGUAGE plpgsql; -SELECT d.windows AS w1, d.chunks AS c1, d.bytes AS b1, d.last_wm AS wm1, - (d.windows > 1 AND d.contiguous AND d.chunks = :base_chunks::bigint - AND d.bytes = :base_bytes::bigint AND d.last_wm = :base_wm::bigint) AS cap1_ok +SELECT d.windows AS w1, d.nrows AS c1, d.maxbytes AS b1, d.last_wm AS wm1, + (d.windows > 1 AND d.contiguous AND d.nrows = :base_rows::bigint + AND d.maxbytes <= 200000 + 3 * 262144 AND d.last_wm = :base_wm::bigint) AS cap1_ok FROM pg_temp.drain_capped(200000) d \gset \if :cap1_ok \echo [PASS] (:testid) a 200 KB budget splits the stream into :w1 windows that tile it exactly \else -\echo [FAIL] (:testid) windows=:w1 chunks=:c1/:base_chunks bytes=:b1/:base_bytes wm=:wm1/:base_wm +\echo [FAIL] (:testid) windows=:w1 rows=:c1/:base_rows maxwindow=:b1 wm=:wm1/:base_wm SELECT (:fail::int + 1) AS fail \gset \endif -SELECT d.windows AS w2, d.chunks AS c2, d.bytes AS b2, d.last_wm AS wm2, - (d.windows > 1 AND d.contiguous AND d.chunks = :base_chunks::bigint - AND d.bytes = :base_bytes::bigint AND d.last_wm = :base_wm::bigint) AS cap2_ok +SELECT d.windows AS w2, d.nrows AS c2, d.maxbytes AS b2, d.last_wm AS wm2, + (d.windows > 1 AND d.contiguous AND d.nrows = :base_rows::bigint + AND d.maxbytes <= 1 + 3 * 262144 AND d.last_wm = :base_wm::bigint) AS cap2_ok FROM pg_temp.drain_capped(1) d \gset \if :cap2_ok \echo [PASS] (:testid) a 1-byte budget still advances and tiles the stream exactly \else -\echo [FAIL] (:testid) windows=:w2 chunks=:c2/:base_chunks bytes=:b2/:base_bytes wm=:wm2/:base_wm +\echo [FAIL] (:testid) windows=:w2 rows=:c2/:base_rows maxwindow=:b2 wm=:wm2/:base_wm SELECT (:fail::int + 1) AS fail \gset \endif diff --git a/test/review_regressions.c b/test/review_regressions.c index c0ffcf2..1395af8 100644 --- a/test/review_regressions.c +++ b/test/review_regressions.c @@ -721,57 +721,78 @@ static void test_payload_window_cap(void) { } CHECK(sql(db, "COMMIT;") == SQLITE_OK); } + // A third shape, and the one that matters most: ~13 rows fill a 256 KB chunk, so + // transactions of 12 then 13 rows keep every chunk boundary one row inside a + // db_version. Waiting for a chunk to end on a boundary never succeeds here, so a + // cap that only checks after a chunk is built never fires and the whole history + // comes out in a single window. + for (int txn = 0; txn < 20; txn++) { + CHECK(sql(db, "BEGIN;") == SQLITE_OK); + for (int i = 0; i < (txn == 0 ? 12 : 13); i++) { + char buf[128]; + snprintf(buf, sizeof(buf), "INSERT INTO t VALUES('r%03d', randomblob(20000));", row++); + CHECK(sql(db, buf) == SQLITE_OK); + } + CHECK(sql(db, "COMMIT;") == SQLITE_OK); + } // Baseline: the whole window in one uncapped scan. sqlite3_stmt *vm = NULL; - int64_t total_bytes = 0, total_chunks = 0, watermark = 0; - CHECK(sqlite3_prepare_v2(db, "SELECT count(*), sum(payload_size), max(watermark_db_version), " + int64_t total_bytes = 0, total_rows = 0, watermark = 0; + CHECK(sqlite3_prepare_v2(db, "SELECT sum(rows), sum(payload_size), " + "max(CASE WHEN is_final THEN watermark_db_version END), " "max(window_capped) FROM cloudsync_payload_chunks(0,NULL,NULL,false)", -1, &vm, NULL) == SQLITE_OK); if (sqlite3_step(vm) == SQLITE_ROW) { - total_chunks = sqlite3_column_int64(vm, 0); + total_rows = sqlite3_column_int64(vm, 0); total_bytes = sqlite3_column_int64(vm, 1); watermark = sqlite3_column_int64(vm, 2); CHECK(sqlite3_column_int(vm, 3) == 0); // nothing is capped without a budget } sqlite3_finalize(vm); - CHECK(total_chunks > 1 && total_bytes > 0 && watermark == 42); + CHECK(total_rows > 0 && total_bytes > 0 && watermark == 62); // Two budgets: one that spans several chunks, and one below a single chunk so the // "always emit one whole db_version" guarantee is what keeps the drain moving. const int64_t budgets[] = {200000, 1}; for (size_t b = 0; b < sizeof(budgets) / sizeof(budgets[0]); ++b) { - int64_t since = 0, sum_bytes = 0, sum_chunks = 0; + int64_t since = 0, sum_rows = 0; int windows = 0; bool capped = true; while (capped && windows < 100) { char q[512]; snprintf(q, sizeof(q), - "SELECT count(*), coalesce(sum(payload_size),0), max(watermark_db_version), " + "SELECT count(*), coalesce(sum(rows),0), coalesce(sum(payload_size),0), " + "max(CASE WHEN is_final THEN watermark_db_version END), " "max(window_capped), min(db_version_min) " "FROM cloudsync_payload_chunks(%lld,NULL,NULL,false,NULL,NULL,NULL,%lld)", (long long)since, (long long)budgets[b]); vm = NULL; CHECK(sqlite3_prepare_v2(db, q, -1, &vm, NULL) == SQLITE_OK); if (sqlite3_step(vm) != SQLITE_ROW || sqlite3_column_int64(vm, 0) == 0) { sqlite3_finalize(vm); break; } - int64_t chunks = sqlite3_column_int64(vm, 0); - int64_t bytes = sqlite3_column_int64(vm, 1); - int64_t wm = sqlite3_column_int64(vm, 2); - capped = sqlite3_column_int(vm, 3) != 0; - int64_t first = sqlite3_column_int64(vm, 4); + int64_t nrows = sqlite3_column_int64(vm, 1); + int64_t bytes = sqlite3_column_int64(vm, 2); + // A window ends at the first db_version boundary at or after the budget, so + // it can overshoot by at most the chunk that crossed it plus the version in + // progress -- never by the whole history. + CHECK(bytes <= budgets[b] + 3 * 262144); + int64_t wm = sqlite3_column_int64(vm, 3); + capped = sqlite3_column_int(vm, 4) != 0; + int64_t first = sqlite3_column_int64(vm, 5); sqlite3_finalize(vm); CHECK(first == since + 1); // windows abut, no gap and no overlap CHECK(wm > since); // always advances, so the drain terminates - sum_chunks += chunks; - sum_bytes += bytes; + sum_rows += nrows; since = wm; windows++; } CHECK(windows > 1); // the budget really did split the stream CHECK(since == watermark); // and the last window reached the end - CHECK(sum_bytes == total_bytes); // tiles the uncapped stream exactly - CHECK(sum_chunks == total_chunks); + // Rows are the invariant, not bytes or chunks: ending a chunk early at a + // db_version boundary repacks the same rows into more chunks, each carrying + // its own header, so a capped drain legitimately moves a few more bytes. + CHECK(sum_rows == total_rows); } CHECK(close_db(db) == SQLITE_OK); From 8661a11667d9c7829d67aa4068a1172171e14c05 Mon Sep 17 00:00:00 2001 From: Andrea Donetti Date: Thu, 24 Sep 2026 12:42:25 -0600 Subject: [PATCH 05/14] fix: count fragment chunks toward the window budget Regression from dd5a95e. Moving the byte accounting into the ordinary chunk builder left out fragment emission, which returns before reaching it: a history of oversized values spent nothing, so the cap never fired and the window was unbounded again. Reproduced in both backends with 256 KB chunks and twenty transactions holding one 300 KB value each -- every budget emitted the whole 6,006,460-byte history with window_capped = 0. Both fragment emitters now add their encoded bytes, in the same unit the ordinary builder uses. Capping still cannot happen mid-value: the cap is evaluated only once a value's last fragment has been emitted. Both suites gain a purely fragmented history, which the existing cases could not have caught -- they contain no value large enough to fragment. The PostgreSQL case keeps its oversized values in the table that is already synced rather than a second one. A table initialised part-way through a session leaves its later single-statement transactions sharing one db_version there (SQLite assigns three for the same sequence), which would give the fragment window no boundary to end on. That looks like a separate defect and is not addressed here. Co-Authored-By: Claude Opus 5 (1M context) --- src/postgresql/cloudsync_postgresql.c | 4 ++ src/sqlite/cloudsync_sqlite.c | 4 ++ test/postgresql/65_payload_window_cap.sql | 34 ++++++++++++++-- test/review_regressions.c | 48 +++++++++++++++++++++++ 4 files changed, 86 insertions(+), 4 deletions(-) diff --git a/src/postgresql/cloudsync_postgresql.c b/src/postgresql/cloudsync_postgresql.c index efb0dd1..6c406b2 100644 --- a/src/postgresql/cloudsync_postgresql.c +++ b/src/postgresql/cloudsync_postgresql.c @@ -1234,6 +1234,10 @@ static bytea *payload_chunks_emit_pg_fragment(PayloadChunksState *st, cloudsync_ if (rc != DBRES_OK) ereport(ERROR, (errcode(cloudsync_error_sqlstate(data)), errmsg("%s", cloudsync_errmsg(data)))); rc = cloudsync_payload_encode_final(payload, data); if (rc != DBRES_OK) ereport(ERROR, (errcode(cloudsync_error_sqlstate(data)), errmsg("%s", cloudsync_errmsg(data)))); + // A fragment chunk spends the window budget like any other. Fragments are emitted + // here rather than by the ordinary builder, so without this a history made of + // oversized values never spends the budget and the cap never fires. + st->window_bytes += (int64)cloudsync_payload_context_bused(payload); int64 blob_size = 0; char *blob = cloudsync_payload_blob(payload, &blob_size, rows); bytea *result = (bytea *)palloc(VARHDRSZ + blob_size); diff --git a/src/sqlite/cloudsync_sqlite.c b/src/sqlite/cloudsync_sqlite.c index 9aea12f..1d510f8 100644 --- a/src/sqlite/cloudsync_sqlite.c +++ b/src/sqlite/cloudsync_sqlite.c @@ -1230,6 +1230,10 @@ static int payload_chunks_emit_fragment(cloudsync_payload_chunks_cursor *c) { if (rc != SQLITE_OK) { cloudsync_memory_free(payload); return rc; } rc = cloudsync_payload_encode_final(payload, data); if (rc != SQLITE_OK) { cloudsync_memory_free(payload); return rc; } + // A fragment chunk spends the window budget like any other. The fragment paths + // return before the ordinary builder's accounting, so without this a history made + // of oversized values never spends the budget and the cap never fires. + c->window_bytes += (int64_t)cloudsync_payload_context_bused(payload); c->payload = cloudsync_payload_blob(payload, &c->payload_size, &c->rows); cloudsync_memory_free(payload); c->dbv_min = sqlite3_column_int64(c->src, 5); diff --git a/test/postgresql/65_payload_window_cap.sql b/test/postgresql/65_payload_window_cap.sql index 7d29479..101ca6e 100644 --- a/test/postgresql/65_payload_window_cap.sql +++ b/test/postgresql/65_payload_window_cap.sql @@ -65,10 +65,10 @@ SELECT (:fail::int + 1) AS fail \gset -- Drain the same stream in capped windows, following watermark_db_version each time, -- and add up what the windows covered. A budget below one chunk also proves the drain -- still advances: a db_version larger than the whole budget must be emitted in full. -CREATE FUNCTION pg_temp.drain_capped(cap bigint) +CREATE FUNCTION pg_temp.drain_capped_from(start_since bigint, cap bigint) RETURNS TABLE (windows int, nrows bigint, maxbytes bigint, last_wm bigint, contiguous boolean) AS $$ DECLARE - since bigint := 0; + since bigint := start_since; w int := 0; c bigint := 0; b bigint := 0; @@ -95,7 +95,7 @@ END $$ LANGUAGE plpgsql; SELECT d.windows AS w1, d.nrows AS c1, d.maxbytes AS b1, d.last_wm AS wm1, (d.windows > 1 AND d.contiguous AND d.nrows = :base_rows::bigint AND d.maxbytes <= 200000 + 3 * 262144 AND d.last_wm = :base_wm::bigint) AS cap1_ok - FROM pg_temp.drain_capped(200000) d \gset + FROM pg_temp.drain_capped_from(0, 200000) d \gset \if :cap1_ok \echo [PASS] (:testid) a 200 KB budget splits the stream into :w1 windows that tile it exactly \else @@ -106,7 +106,7 @@ SELECT (:fail::int + 1) AS fail \gset SELECT d.windows AS w2, d.nrows AS c2, d.maxbytes AS b2, d.last_wm AS wm2, (d.windows > 1 AND d.contiguous AND d.nrows = :base_rows::bigint AND d.maxbytes <= 1 + 3 * 262144 AND d.last_wm = :base_wm::bigint) AS cap2_ok - FROM pg_temp.drain_capped(1) d \gset + FROM pg_temp.drain_capped_from(0, 1) d \gset \if :cap2_ok \echo [PASS] (:testid) a 1-byte budget still advances and tiles the stream exactly \else @@ -125,6 +125,32 @@ SELECT count(*) AS res_chunks, min(db_version_min) AS res_first, SELECT (:fail::int + 1) AS fail \gset \endif +-- A history of oversized values is emitted entirely as fragment chunks, which the +-- ordinary chunk builder never produces. Those bytes still have to spend the budget, or +-- such a history never reaches the cap at all. +INSERT INTO items (id, v) SELECT 'f1', (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 18750) g); +INSERT INTO items (id, v) SELECT 'f2', (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 18750) g); +INSERT INTO items (id, v) SELECT 'f3', (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 18750) g); +INSERT INTO items (id, v) SELECT 'f4', (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 18750) g); +INSERT INTO items (id, v) SELECT 'f5', (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 18750) g); +INSERT INTO items (id, v) SELECT 'f6', (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 18750) g); +INSERT INTO items (id, v) SELECT 'f7', (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 18750) g); +INSERT INTO items (id, v) SELECT 'f8', (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 18750) g); + +SELECT sum(rows) AS frag_rows, max(watermark_db_version) AS frag_wm + FROM cloudsync_payload_chunks(:base_wm, NULL, NULL, false) \gset + +SELECT d.windows AS w3, d.nrows AS c3, d.last_wm AS wm3, + (d.windows > 1 AND d.contiguous AND d.nrows = :frag_rows::bigint + AND d.last_wm = :frag_wm::bigint) AS cap3_ok + FROM pg_temp.drain_capped_from(:base_wm, 200000) d \gset +\if :cap3_ok +\echo [PASS] (:testid) a purely fragmented history is capped into :w3 windows +\else +\echo [FAIL] (:testid) windows=:w3 rows=:c3/:frag_rows wm=:wm3/:frag_wm +SELECT (:fail::int + 1) AS fail \gset +\endif + \connect postgres \ir helper_psql_conn_setup.sql DROP DATABASE IF EXISTS cloudsync_test_65; diff --git a/test/review_regressions.c b/test/review_regressions.c index 1395af8..4ede5d1 100644 --- a/test/review_regressions.c +++ b/test/review_regressions.c @@ -796,6 +796,54 @@ static void test_payload_window_cap(void) { } CHECK(close_db(db) == SQLITE_OK); + + // A history of oversized values is emitted entirely as fragment chunks, which the + // ordinary chunk builder never produces. Those bytes still have to spend the budget, + // or such a history never reaches the cap at all. + db = open_db(); + CHECK(sql(db, "CREATE TABLE t(id TEXT PRIMARY KEY NOT NULL, v BLOB);" + "SELECT cloudsync_init('t');" + "SELECT cloudsync_set('payload_max_chunk_size','262144');") == SQLITE_OK); + for (int i = 0; i < 20; i++) { + char buf[128]; + snprintf(buf, sizeof(buf), "INSERT INTO t VALUES('f%03d', randomblob(300000));", i); + CHECK(sql(db, buf) == SQLITE_OK); + } + int64_t frag_rows = 0, frag_wm = 0; + vm = NULL; + CHECK(sqlite3_prepare_v2(db, "SELECT sum(rows), max(CASE WHEN is_final THEN watermark_db_version END) " + "FROM cloudsync_payload_chunks(0,NULL,NULL,false)", -1, &vm, NULL) == SQLITE_OK); + if (sqlite3_step(vm) == SQLITE_ROW) { frag_rows = sqlite3_column_int64(vm, 0); frag_wm = sqlite3_column_int64(vm, 1); } + sqlite3_finalize(vm); + CHECK(frag_rows > 0 && frag_wm == 20); + + int64_t since = 0, seen = 0; + int windows = 0; + bool capped = true; + while (capped && windows < 100) { + vm = NULL; + char q[384]; + snprintf(q, sizeof(q), + "SELECT count(*), coalesce(sum(rows),0), " + "max(CASE WHEN is_final THEN watermark_db_version END), max(window_capped), " + "min(db_version_min) " + "FROM cloudsync_payload_chunks(%lld,NULL,NULL,false,NULL,NULL,NULL,200000)", + (long long)since); + CHECK(sqlite3_prepare_v2(db, q, -1, &vm, NULL) == SQLITE_OK); + if (sqlite3_step(vm) != SQLITE_ROW || sqlite3_column_int64(vm, 0) == 0) { sqlite3_finalize(vm); break; } + seen += sqlite3_column_int64(vm, 1); + int64_t wm = sqlite3_column_int64(vm, 2); + capped = sqlite3_column_int(vm, 3) != 0; + CHECK(sqlite3_column_int64(vm, 4) == since + 1); + CHECK(wm > since); + sqlite3_finalize(vm); + since = wm; + windows++; + } + CHECK(windows > 1); // the budget split a purely fragmented history + CHECK(since == frag_wm); + CHECK(seen == frag_rows); + CHECK(close_db(db) == SQLITE_OK); } int main(void) { From 98081522eeb3f5d12e1ddd9dc73ca3d4974035f5 Mon Sep 17 00:00:00 2001 From: Andrea Donetti Date: Thu, 24 Sep 2026 12:47:53 -0600 Subject: [PATCH 06/14] test(postgres): cite issue #69 for the db_version collapse the test works around The fragment case keeps its oversized values in the already-synced table because by that point the session has held snapshots open, and #69 defect 2 leaves the cached db_version unreloaded until txid_snapshot_xmin changes -- fresh single-statement transactions then collapse onto one db_version, giving the window no boundary to end on. Co-Authored-By: Claude Opus 5 (1M context) --- test/postgresql/65_payload_window_cap.sql | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/test/postgresql/65_payload_window_cap.sql b/test/postgresql/65_payload_window_cap.sql index 101ca6e..63fdcb2 100644 --- a/test/postgresql/65_payload_window_cap.sql +++ b/test/postgresql/65_payload_window_cap.sql @@ -128,6 +128,12 @@ SELECT (:fail::int + 1) AS fail \gset -- A history of oversized values is emitted entirely as fragment chunks, which the -- ordinary chunk builder never produces. Those bytes still have to spend the budget, or -- such a history never reaches the cap at all. +-- +-- These rows go in the table already synced above rather than a second one. By this +-- point the session has held snapshots open (the drains above), and issue #69 defect 2 +-- means the cached db_version is only reloaded when txid_snapshot_xmin changes -- so +-- fresh single-statement transactions collapse onto one db_version and the window would +-- have no boundary to end on. Keeping to one table avoids that while it is open. INSERT INTO items (id, v) SELECT 'f1', (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 18750) g); INSERT INTO items (id, v) SELECT 'f2', (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 18750) g); INSERT INTO items (id, v) SELECT 'f3', (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 18750) g); From fbc5947376e1bf31d901e11424b03507308a4e75 Mon Sep 17 00:00:00 2001 From: Andrea Donetti Date: Thu, 24 Sep 2026 16:23:08 -0600 Subject: [PATCH 07/14] test(postgres): correct why the fragment case stays in one table The comment blamed the pinned-xmin cause in issue #69. The actual reason is the issue step 1 defect: the db_version reload reads only the first synced table s maximum, so a second table s writes collapse onto one db_version, leaving the window no boundary to end on. Writes to the first table keep advancing in the same session, which is what the earlier attribution could not explain. Co-Authored-By: Claude Opus 5 (1M context) --- test/postgresql/65_payload_window_cap.sql | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/test/postgresql/65_payload_window_cap.sql b/test/postgresql/65_payload_window_cap.sql index 63fdcb2..9457f5d 100644 --- a/test/postgresql/65_payload_window_cap.sql +++ b/test/postgresql/65_payload_window_cap.sql @@ -129,11 +129,11 @@ SELECT (:fail::int + 1) AS fail \gset -- ordinary chunk builder never produces. Those bytes still have to spend the budget, or -- such a history never reaches the cap at all. -- --- These rows go in the table already synced above rather than a second one. By this --- point the session has held snapshots open (the drains above), and issue #69 defect 2 --- means the cached db_version is only reloaded when txid_snapshot_xmin changes -- so --- fresh single-statement transactions collapse onto one db_version and the window would --- have no boundary to end on. Keeping to one table avoids that while it is open. +-- These rows go in the table already synced above rather than a second one: on +-- PostgreSQL the db_version reload reads only the first synced table's maximum, so a +-- second table's writes collapse onto a single db_version and the window would have no +-- boundary to end on. Tracked as the step 1 defect of issue #69; keeping to one table +-- avoids it while that is open. INSERT INTO items (id, v) SELECT 'f1', (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 18750) g); INSERT INTO items (id, v) SELECT 'f2', (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 18750) g); INSERT INTO items (id, v) SELECT 'f3', (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 18750) g); From 4e8056699dfd580b6aa18cf0536380c55d82143f Mon Sep 17 00:00:00 2001 From: Andrea Donetti Date: Thu, 24 Sep 2026 17:02:49 -0600 Subject: [PATCH 08/14] test(postgres): give the fragment case its own table again The oversized values were kept in the already-synced table to avoid the issue #69 step 1 defect, where the db_version reload read only the first synced table's maximum and a second table's writes collapsed onto one db_version, leaving the window no boundary to end on. #78 fixed that, so the case goes back to a table of its own, which is what it was meant to be. It now also exercises the fix: eight single-statement transactions into a second table have to take eight db_versions for the fragmented history to be capped into eight windows. Co-Authored-By: Claude Opus 5 (1M context) --- test/postgresql/65_payload_window_cap.sql | 25 ++++++++++------------- 1 file changed, 11 insertions(+), 14 deletions(-) diff --git a/test/postgresql/65_payload_window_cap.sql b/test/postgresql/65_payload_window_cap.sql index 9457f5d..02b9822 100644 --- a/test/postgresql/65_payload_window_cap.sql +++ b/test/postgresql/65_payload_window_cap.sql @@ -21,6 +21,8 @@ CREATE DATABASE cloudsync_test_65; CREATE EXTENSION IF NOT EXISTS cloudsync; CREATE TABLE items (id TEXT PRIMARY KEY NOT NULL, v BYTEA); SELECT cloudsync_init('items', 'CLS', 1) AS _init \gset +CREATE TABLE big (id TEXT PRIMARY KEY NOT NULL, v BYTEA); +SELECT cloudsync_init('big', 'CLS', 1) AS _init2 \gset -- The smallest chunk size the setting allows, so the rows below span several chunks. SELECT cloudsync_set('payload_max_chunk_size', '262144') AS _chunk \gset @@ -128,20 +130,15 @@ SELECT (:fail::int + 1) AS fail \gset -- A history of oversized values is emitted entirely as fragment chunks, which the -- ordinary chunk builder never produces. Those bytes still have to spend the budget, or -- such a history never reaches the cap at all. --- --- These rows go in the table already synced above rather than a second one: on --- PostgreSQL the db_version reload reads only the first synced table's maximum, so a --- second table's writes collapse onto a single db_version and the window would have no --- boundary to end on. Tracked as the step 1 defect of issue #69; keeping to one table --- avoids it while that is open. -INSERT INTO items (id, v) SELECT 'f1', (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 18750) g); -INSERT INTO items (id, v) SELECT 'f2', (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 18750) g); -INSERT INTO items (id, v) SELECT 'f3', (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 18750) g); -INSERT INTO items (id, v) SELECT 'f4', (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 18750) g); -INSERT INTO items (id, v) SELECT 'f5', (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 18750) g); -INSERT INTO items (id, v) SELECT 'f6', (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 18750) g); -INSERT INTO items (id, v) SELECT 'f7', (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 18750) g); -INSERT INTO items (id, v) SELECT 'f8', (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 18750) g); + +INSERT INTO big (id, v) SELECT 'f1', (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 18750) g); +INSERT INTO big (id, v) SELECT 'f2', (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 18750) g); +INSERT INTO big (id, v) SELECT 'f3', (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 18750) g); +INSERT INTO big (id, v) SELECT 'f4', (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 18750) g); +INSERT INTO big (id, v) SELECT 'f5', (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 18750) g); +INSERT INTO big (id, v) SELECT 'f6', (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 18750) g); +INSERT INTO big (id, v) SELECT 'f7', (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 18750) g); +INSERT INTO big (id, v) SELECT 'f8', (SELECT decode(string_agg(md5(random()::text || g::text), ''), 'hex') FROM generate_series(1, 18750) g); SELECT sum(rows) AS frag_rows, max(watermark_db_version) AS frag_wm FROM cloudsync_payload_chunks(:base_wm, NULL, NULL, false) \gset From 05b3e4915cfa8c61cd1a9644efb7cc8be2378475 Mon Sep 17 00:00:00 2001 From: Andrea Donetti Date: Fri, 25 Sep 2026 10:18:32 -0600 Subject: [PATCH 09/14] fix: carry the window budget across a stream fetched one chunk per call The cap counted only the bytes of a single scan. A stateless caller fetches one chunk per call and resumes through resume_*, so every call started the count at zero and the budget was never reached: against any budget larger than one chunk the whole window came out uncapped, however long the drain ran. That is how the /check job pages, so as it stood the cap did nothing for the only caller that needs it. Carry the spend the same way the stream position is already carried. A new window_bytes output reports what the window has spent so far, and a new 9th argument resume_window_bytes seeds it. The server passes the same max_window_bytes on every call and hands window_bytes back. One unit, no arithmetic for the caller to get wrong, and symmetric with next_*/resume_*. The alternative -- having the caller subtract what it has spent -- was rejected twice over: the extension counts encoded bytes before compression while a caller only sees the compressed payload_size, so the two budgets would diverge arbitrarily on compressible data; and a remaining budget reaching zero reads as "no cap", disabling the protection exactly when it is due. Both new inputs stay last, so positional arguments 1..7 keep their meaning, and window_bytes/window_capped are outputs, which do not move them either. An 8-argument caller still works: resume_window_bytes defaults to 0, which is the previous per-call behaviour. Both suites gain a drain that pages one chunk per call under a budget spanning several chunks. The existing tests drain a whole window in one scan, where the count accumulates naturally, and so could not have caught this. Co-Authored-By: Claude Opus 5 (1M context) --- API.md | 18 +++-- CHANGELOG.md | 2 +- src/postgresql/cloudsync.sql.in | 6 +- src/postgresql/cloudsync_postgresql.c | 10 ++- .../migrations/cloudsync--1.1--1.2.sql | 14 ++-- src/sqlite/cloudsync_sqlite.c | 31 +++++-- test/postgresql/65_payload_window_cap.sql | 57 +++++++++++++ test/review_regressions.c | 81 +++++++++++++++++++ 8 files changed, 194 insertions(+), 25 deletions(-) diff --git a/API.md b/API.md index b08c326..72d8c79 100644 --- a/API.md +++ b/API.md @@ -30,7 +30,7 @@ This document provides a reference for the SQL functions provided by the `sqlite - [Payload Functions](#payload-functions) - [`cloudsync_payload_encode()`](#cloudsync_payload_encodetbl-pk-col_name-col_value-col_version-db_version-site_id-cl-seq) - [`cloudsync_payload_blob_checked()`](#cloudsync_payload_blob_checkedsince_db_version-since_seq-filter_site_id-exclude_filter_site_id-max_estimated_payload_size) - - [`cloudsync_payload_chunks()`](#cloudsync_payload_chunkssince_db_version-filter_site_id-until_db_version-exclude_filter_site_id-resume_db_version-resume_seq-resume_frag_offset-max_window_bytes) + - [`cloudsync_payload_chunks()`](#cloudsync_payload_chunkssince_db_version-filter_site_id-until_db_version-exclude_filter_site_id-resume_db_version-resume_seq-resume_frag_offset-max_window_bytes-resume_window_bytes) - [`cloudsync_payload_apply()`](#cloudsync_payload_applypayload) - [Network Functions](#network-functions) - [`cloudsync_network_init()`](#cloudsync_network_initmanageddatabaseid) @@ -57,7 +57,7 @@ The following settings are supported: | Key | Description | Default | Minimum | Maximum | |---|---|---:|---:|---:| -| `payload_max_chunk_size` | Maximum transport payload size generated by [`cloudsync_payload_chunks()`](#cloudsync_payload_chunkssince_db_version-filter_site_id-until_db_version-exclude_filter_site_id-resume_db_version-resume_seq-resume_frag_offset-max_window_bytes). Values outside the range are clamped. | `5242880` (5 MB) | `262144` (256 KB) | `33554432` (32 MB) | +| `payload_max_chunk_size` | Maximum transport payload size generated by [`cloudsync_payload_chunks()`](#cloudsync_payload_chunkssince_db_version-filter_site_id-until_db_version-exclude_filter_site_id-resume_db_version-resume_seq-resume_frag_offset-max_window_bytes-resume_window_bytes). Values outside the range are clamped. | `5242880` (5 MB) | `262144` (256 KB) | `33554432` (32 MB) | | `network_connect_timeout` | Seconds allowed to connect, name resolution included. Never longer than the request's own deadline. | `30` | `1` | `86400` | | `network_request_timeout` | Seconds allowed for a whole API request (check, upload URL, apply, status). | `300` | `1` | `86400` | | `network_artifact_timeout` | Absolute limit, in seconds, for a payload upload or download. | `3600` | `1` | `86400` | @@ -436,7 +436,7 @@ SELECT cloudsync_uuid_text(cloudsync_siteid(), false); -- 0190a1b2c3d47e5f8a9b ### `cloudsync_uuid_blob(uuid)` -**Description:** Converts a UUID string into its 16-byte binary form. This is the inverse of [`cloudsync_uuid_text()`](#cloudsync_uuid_textuuid-dash_format) and lets string-based callers (for example, an HTTP `/check` endpoint holding a stringified `site_id`) pass a `site_id` to [`cloudsync_payload_chunks()`](#cloudsync_payload_chunkssince_db_version-filter_site_id-until_db_version-exclude_filter_site_id-resume_db_version-resume_seq-resume_frag_offset-max_window_bytes). +**Description:** Converts a UUID string into its 16-byte binary form. This is the inverse of [`cloudsync_uuid_text()`](#cloudsync_uuid_textuuid-dash_format) and lets string-based callers (for example, an HTTP `/check` endpoint holding a stringified `site_id`) pass a `site_id` to [`cloudsync_payload_chunks()`](#cloudsync_payload_chunkssince_db_version-filter_site_id-until_db_version-exclude_filter_site_id-resume_db_version-resume_seq-resume_frag_offset-max_window_bytes-resume_window_bytes). **Parameters:** @@ -505,7 +505,7 @@ SELECT cloudsync_commit_alter('my_table'); **Description:** Encodes rows from `cloudsync_changes` into a single monolithic payload. This is the legacy payload API and remains fully supported for backward compatibility. -Use this API when the expected payload size is modest or when you need to interoperate with callers that expect a single BLOB. For large rowsets or large individual BLOB/TEXT values, prefer [`cloudsync_payload_chunks()`](#cloudsync_payload_chunkssince_db_version-filter_site_id-until_db_version-exclude_filter_site_id-resume_db_version-resume_seq-resume_frag_offset-max_window_bytes), which splits transport payloads according to `payload_max_chunk_size`. +Use this API when the expected payload size is modest or when you need to interoperate with callers that expect a single BLOB. For large rowsets or large individual BLOB/TEXT values, prefer [`cloudsync_payload_chunks()`](#cloudsync_payload_chunkssince_db_version-filter_site_id-until_db_version-exclude_filter_site_id-resume_db_version-resume_seq-resume_frag_offset-max_window_bytes-resume_window_bytes), which splits transport payloads according to `payload_max_chunk_size`. **Parameters:** The function is an aggregate over the columns returned by `cloudsync_changes`: @@ -567,7 +567,7 @@ SELECT cloudsync_payload_blob_checked( --- -### `cloudsync_payload_chunks([since_db_version], [filter_site_id], [until_db_version], [exclude_filter_site_id], [resume_db_version], [resume_seq], [resume_frag_offset], [max_window_bytes])` +### `cloudsync_payload_chunks([since_db_version], [filter_site_id], [until_db_version], [exclude_filter_site_id], [resume_db_version], [resume_seq], [resume_frag_offset], [max_window_bytes]), [resume_window_bytes])` **Description:** Generates sync payloads as a stream of transport-sized chunks. It is the chunk-aware evolution of [`cloudsync_payload_encode()`](#cloudsync_payload_encodetbl-pk-col_name-col_value-col_version-db_version-site_id-cl-seq), designed for large rowsets and for single BLOB/TEXT values that are larger than the configured chunk size. @@ -592,7 +592,10 @@ When a single encoded column value does not fit in one chunk, CloudSync transpar - `resume_db_version`, `resume_seq`, `resume_frag_offset` (INTEGER/BIGINT, optional): Resume a stream where a previous call stopped, instead of replaying it from `since_db_version`. Pass back the `next_db_version`, `next_seq` and `next_frag_offset` of the last chunk received, together with the `watermark_db_version` as `until_db_version` so the window stays fixed. This lets a stateless caller page one chunk per request without holding a cursor, and the resume is a seek rather than a re-scan. Supplying `resume_db_version` replaces `since_db_version` as the lower bound; leave all three unset (or `NULL`) for a fresh stream. - `max_window_bytes` (INTEGER/BIGINT, optional): Cap the whole stream at roughly this many payload bytes, instead of preparing the entire window. When the budget is spent the stream ends at the next complete `db_version`, reporting `is_final` with `watermark_db_version` lowered to that point and `window_capped` set. The result is an ordinary *complete* stream over a smaller window: checkpoint at the reported watermark and call again with that value as `since_db_version` to continue. Values of `0` or `NULL` mean no cap, which is the default. + A caller that fetches the whole window in one query can stop here. A caller that fetches **one chunk per query** and resumes through `resume_*` must also carry `resume_window_bytes` (below), otherwise each query counts only its own chunk, the budget is never reached, and nothing is ever capped. + Use it to bound how long one preparation takes on a large history, and note three properties. A window always ends on a `db_version` boundary, so a transaction larger than the budget is still emitted whole and the budget is an approximate floor rather than a hard ceiling. Reaching that boundary can mean ending a chunk before it is full, so a capped drain packs the same changes into slightly more chunks than an uncapped one. And because a window never ends empty, repeated calls always make progress. +- `resume_window_bytes` (INTEGER/BIGINT, optional): Budget already spent by earlier queries of this window. Pass back the `window_bytes` the previous query returned, exactly as `resume_db_version` carries the stream position. Omit it (or pass `0`) to start a window; it has no effect without `max_window_bytes`. **Returns:** A rowset with one row per chunk: @@ -608,6 +611,7 @@ When a single encoded column value does not fit in one chunk, CloudSync transpar | `next_db_version`, `next_seq`, `next_frag_offset` | Resume point immediately after this chunk. Pass them back as the `resume_*` arguments to continue the stream. | | `is_final` | True on the last chunk of the stream. A capped window is a complete stream, so its last chunk reports `is_final` too. | | `window_capped` | True when the stream ended because `max_window_bytes` was reached rather than because the window was drained — i.e. more changes exist past `watermark_db_version`. | +| `window_bytes` | Budget spent by this window so far, including earlier queries that passed `resume_window_bytes`. Pass it back as `resume_window_bytes` on the next query of the same window. | **SQLite usage:** `cloudsync_payload_chunks` is exposed as a virtual table with hidden constraint columns: @@ -667,7 +671,7 @@ On PostgreSQL, apply chunks as individual statements from the transport/client l - Legacy payloads generated by older SQLite Sync versions. - Monolithic payloads generated by [`cloudsync_payload_encode()`](#cloudsync_payload_encodetbl-pk-col_name-col_value-col_version-db_version-site_id-cl-seq). -- Chunk-fragment payloads generated by [`cloudsync_payload_chunks()`](#cloudsync_payload_chunkssince_db_version-filter_site_id-until_db_version-exclude_filter_site_id-resume_db_version-resume_seq-resume_frag_offset-max_window_bytes). +- Chunk-fragment payloads generated by [`cloudsync_payload_chunks()`](#cloudsync_payload_chunkssince_db_version-filter_site_id-until_db_version-exclude_filter_site_id-resume_db_version-resume_seq-resume_frag_offset-max_window_bytes-resume_window_bytes). When a v3 fragment payload is received, CloudSync stores the fragment in an internal table and returns after applying zero or more completed values. Once the final fragment for a value is received, the completed value is validated and applied. Fragments can arrive in any order, and duplicate fragment delivery is idempotent. Applying a fragment never moves the receive checkpoint. On PostgreSQL, pieces of one value applied by concurrent transactions wait for each other under `READ COMMITTED`, and fail with a retryable serialization error under `SERIALIZABLE` when they conflict; a fragment is refused under `REPEATABLE READ`, where a transaction could miss a piece committed while it waited. @@ -786,7 +790,7 @@ This means: if you get JSON back, the server was reachable and the network proto **Description:** Sends all unsent local changes to the remote server. -The send path streams payloads through [`cloudsync_payload_chunks()`](#cloudsync_payload_chunkssince_db_version-filter_site_id-until_db_version-exclude_filter_site_id-resume_db_version-resume_seq-resume_frag_offset-max_window_bytes), so `payload_max_chunk_size` also limits the payloads generated for network transport. Each generated chunk is uploaded/applied independently; the local send checkpoint is advanced only after the chunk stream completes successfully. +The send path streams payloads through [`cloudsync_payload_chunks()`](#cloudsync_payload_chunkssince_db_version-filter_site_id-until_db_version-exclude_filter_site_id-resume_db_version-resume_seq-resume_frag_offset-max_window_bytes-resume_window_bytes), so `payload_max_chunk_size` also limits the payloads generated for network transport. Each generated chunk is uploaded/applied independently; the local send checkpoint is advanced only after the chunk stream completes successfully. 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. diff --git a/CHANGELOG.md b/CHANGELOG.md index 435ec96..ee59546 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,7 +9,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/). ### 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. -- **`cloudsync_payload_chunks()` can bound how much one call prepares.** Preparing a large history is unbounded work: a big enough tenant cannot finish inside a caller's time budget, and because nothing is durable until the stream reports `is_final`, an attempt that runs out of time keeps no progress and the next one restarts from the first chunk. The new `max_window_bytes` argument ends the stream once roughly that many payload bytes have been emitted, at the next complete database version. What comes back is an ordinary *complete* stream over a smaller window — `is_final` with `watermark_db_version` lowered to that point — so the caller checkpoints there and calls again to continue, with no resumable state to keep anywhere. A new `window_capped` output says the stream stopped on the budget rather than because the window was drained, so more changes exist past the watermark. Two properties are worth knowing: a window always ends on a database version boundary, so a transaction larger than the budget is still emitted whole and the budget is an approximate floor rather than a hard ceiling; and a window never ends empty, so repeated calls always make progress. Reaching that boundary can require ending a chunk before it is full, so a capped drain packs the same changes into slightly more chunks than an uncapped one. Unset, the function behaves exactly as before. +- **`cloudsync_payload_chunks()` can bound how much one call prepares.** Preparing a large history is unbounded work: a big enough tenant cannot finish inside a caller's time budget, and because nothing is durable until the stream reports `is_final`, an attempt that runs out of time keeps no progress and the next one restarts from the first chunk. The new `max_window_bytes` argument ends the stream once roughly that many payload bytes have been emitted, at the next complete database version. What comes back is an ordinary *complete* stream over a smaller window — `is_final` with `watermark_db_version` lowered to that point — so the caller checkpoints there and calls again to continue, with no resumable state to keep anywhere. A new `window_capped` output says the stream stopped on the budget rather than because the window was drained, so more changes exist past the watermark. A caller that fetches one chunk per query and resumes through `resume_*` also passes the `window_bytes` output back as `resume_window_bytes`, which carries the budget spent across those queries exactly as `resume_db_version` carries the stream position; without it each query would count only its own chunk and the budget would never be reached. Two properties are worth knowing: a window always ends on a database version boundary, so a transaction larger than the budget is still emitted whole and the budget is an approximate floor rather than a hard ceiling; and a window never ends empty, so repeated calls always make progress. Reaching that boundary can require ending a chunk before it is full, so a capped drain packs the same changes into slightly more chunks than an uncapped one. Unset, the function behaves exactly as before. ### Changed diff --git a/src/postgresql/cloudsync.sql.in b/src/postgresql/cloudsync.sql.in index 0cc4594..9df66b7 100644 --- a/src/postgresql/cloudsync.sql.in +++ b/src/postgresql/cloudsync.sql.in @@ -157,7 +157,8 @@ CREATE OR REPLACE FUNCTION cloudsync_payload_chunks( resume_db_version bigint DEFAULT NULL, resume_seq bigint DEFAULT NULL, resume_frag_offset bigint DEFAULT NULL, - max_window_bytes bigint DEFAULT NULL + max_window_bytes bigint DEFAULT NULL, + resume_window_bytes bigint DEFAULT NULL ) RETURNS TABLE ( payload bytea, @@ -171,7 +172,8 @@ RETURNS TABLE ( next_seq bigint, next_frag_offset bigint, is_final boolean, - window_capped boolean + window_capped boolean, + window_bytes bigint ) AS 'MODULE_PATHNAME', 'cloudsync_payload_chunks' LANGUAGE C VOLATILE; diff --git a/src/postgresql/cloudsync_postgresql.c b/src/postgresql/cloudsync_postgresql.c index 6c406b2..b95343d 100644 --- a/src/postgresql/cloudsync_postgresql.c +++ b/src/postgresql/cloudsync_postgresql.c @@ -1403,6 +1403,11 @@ Datum cloudsync_payload_chunks(PG_FUNCTION_ARGS) { // Cap the whole prepared window, not one chunk: <= 0 and NULL both mean no cap. int64 window_cap = PG_ARGISNULL(7) ? 0 : PG_GETARG_INT64(7); st->max_window_bytes = (window_cap > 0) ? window_cap : 0; + // Budget already spent by earlier calls of this window. State is created fresh + // per call, so a caller paging one chunk at a time seeds it here; a caller + // draining the window in one call leaves it at 0 and it accumulates. + int64 window_spent = PG_ARGISNULL(8) ? 0 : PG_GETARG_INT64(8); + st->window_bytes = (window_spent > 0) ? window_spent : 0; // Site filter resolution: // exclude=true -> all sites except filter_site_id (CHECK path); site required // filter given -> only that site @@ -1542,8 +1547,8 @@ Datum cloudsync_payload_chunks(PG_FUNCTION_ARGS) { } } - Datum outvals[12]; - bool outnulls[12] = {false,false,false,false,false,false,false,false,false,false,false,false}; + Datum outvals[13]; + bool outnulls[13] = {false,false,false,false,false,false,false,false,false,false,false,false,false}; outvals[0] = PointerGetDatum(payload); outvals[1] = Int64GetDatum(st->chunk_index++); outvals[2] = Int64GetDatum(VARSIZE_ANY_EXHDR(payload)); @@ -1556,6 +1561,7 @@ Datum cloudsync_payload_chunks(PG_FUNCTION_ARGS) { outvals[9] = Int64GetDatum(next_frag); outvals[10] = BoolGetDatum(is_final); outvals[11] = BoolGetDatum(st->window_capped); + outvals[12] = Int64GetDatum(st->window_bytes); HeapTuple outtup = heap_form_tuple(st->outdesc, outvals, outnulls); SRF_RETURN_NEXT(funcctx, HeapTupleGetDatum(outtup)); } diff --git a/src/postgresql/migrations/cloudsync--1.1--1.2.sql b/src/postgresql/migrations/cloudsync--1.1--1.2.sql index 45a3f00..fa280b0 100644 --- a/src/postgresql/migrations/cloudsync--1.1--1.2.sql +++ b/src/postgresql/migrations/cloudsync--1.1--1.2.sql @@ -1,10 +1,12 @@ -- CloudSync PostgreSQL extension upgrade: 1.1 -> 1.2 -- -- Bounds what one call to cloudsync_payload_chunks() prepares: --- * new max_window_bytes input, declared last so the existing positional --- arguments 1..7 keep their meaning +-- * new max_window_bytes and resume_window_bytes inputs, declared last so the +-- existing positional arguments 1..7 keep their meaning -- * new window_capped output, true when the scan stopped on the budget --- rather than because the window was drained +-- rather than because the window was drained, and window_bytes, the budget +-- spent so far -- pass it back as resume_window_bytes to carry the budget +-- across a stream fetched one chunk per call -- -- Both are optional: without max_window_bytes the function behaves exactly as -- it did in 1.1. @@ -26,7 +28,8 @@ CREATE OR REPLACE FUNCTION cloudsync_payload_chunks( resume_db_version bigint DEFAULT NULL, resume_seq bigint DEFAULT NULL, resume_frag_offset bigint DEFAULT NULL, - max_window_bytes bigint DEFAULT NULL + max_window_bytes bigint DEFAULT NULL, + resume_window_bytes bigint DEFAULT NULL ) RETURNS TABLE ( payload bytea, @@ -40,7 +43,8 @@ RETURNS TABLE ( next_seq bigint, next_frag_offset bigint, is_final boolean, - window_capped boolean + window_capped boolean, + window_bytes bigint ) AS 'MODULE_PATHNAME', 'cloudsync_payload_chunks' LANGUAGE C VOLATILE; diff --git a/src/sqlite/cloudsync_sqlite.c b/src/sqlite/cloudsync_sqlite.c index 1d510f8..3e6e236 100644 --- a/src/sqlite/cloudsync_sqlite.c +++ b/src/sqlite/cloudsync_sqlite.c @@ -1054,12 +1054,18 @@ static int payload_chunks_connect(sqlite3 *db, void *aux, int argc, const char * // back as the resume_* inputs (cols 15..17) to continue the drain without // a spool table — O(1) seek per chunk instead of replaying from since. "next_db_version INTEGER, next_seq INTEGER, next_frag_offset INTEGER, is_final INTEGER, " - // window_capped (col 15) is an output, so it does not shift the hidden columns - // a table-valued call binds positionally. max_window_bytes is declared last for - // the same reason: it becomes argument 8, leaving arguments 1..7 as they were. - "window_capped INTEGER, " + // window_capped and window_bytes (cols 15..16) are outputs, so they do not shift + // the hidden columns a table-valued call binds positionally. The two budget + // inputs are declared last for the same reason: they become arguments 8 and 9, + // leaving arguments 1..7 as they were. + // + // window_bytes/resume_window_bytes carry the budget spent across calls exactly + // as next_*/resume_* carry the stream position. Without it a caller that fetches + // one chunk per call -- which is how a stateless /check pages -- restarts the + // count every time and never reaches the budget at all. + "window_capped INTEGER, window_bytes INTEGER, " "resume_db_version HIDDEN, resume_seq HIDDEN, resume_frag_offset HIDDEN, " - "max_window_bytes HIDDEN)"); + "max_window_bytes HIDDEN, resume_window_bytes HIDDEN)"); if (rc != SQLITE_OK) return rc; cloudsync_payload_chunks_vtab *p = sqlite3_malloc64(sizeof(*p)); if (!p) return SQLITE_NOMEM; @@ -1097,9 +1103,10 @@ static int payload_chunks_best_index(sqlite3_vtab *vtab, sqlite3_index_info *idx // in a fixed order regardless of how SQLite presents constraints. idxNum bit k // is set when handled_cols[k] is bound; xFilter reads argv in this same order. // bit0=since_db_version(7) bit1=site_id(8) bit2=until_db_version(9) - // bit3=exclude_filter_site_id(10) bit4=resume_db_version(16) - // bit5=resume_seq(17) bit6=resume_frag_offset(18) bit7=max_window_bytes(19) - static const int handled_cols[] = {7, 8, 9, 10, 16, 17, 18, 19}; + // bit3=exclude_filter_site_id(10) bit4=resume_db_version(17) + // bit5=resume_seq(18) bit6=resume_frag_offset(19) bit7=max_window_bytes(20) + // bit8=resume_window_bytes(21) + static const int handled_cols[] = {7, 8, 9, 10, 17, 18, 19, 20, 21}; int argv_index = 1; int idxnum = 0; for (size_t k = 0; k < sizeof(handled_cols) / sizeof(handled_cols[0]); ++k) { @@ -1422,6 +1429,13 @@ static int payload_chunks_filter(sqlite3_vtab_cursor *cursor, int idxnum, const int64_t cap = sqlite3_value_int64(argv[argi++]); c->max_window_bytes = (cap > 0) ? cap : 0; // <= 0 means no cap } + // Budget already spent by earlier calls of this window. The per-scan reset above + // cleared window_bytes, so a caller paging one chunk at a time seeds it here; a + // caller draining the window in one scan leaves it at 0 and it accumulates. + if (idxnum & 256) { + int64_t spent = sqlite3_value_int64(argv[argi++]); + c->window_bytes = (spent > 0) ? spent : 0; + } // Resolve the site filter: // exclude=true -> all sites except filter_site_id (CHECK path); site required @@ -1537,6 +1551,7 @@ static int payload_chunks_column(sqlite3_vtab_cursor *cursor, sqlite3_context *c case 13: sqlite3_result_int64(ctx, c->next_frag_offset); break; case 14: sqlite3_result_int(ctx, c->is_final ? 1 : 0); break; case 15: sqlite3_result_int(ctx, c->window_capped ? 1 : 0); break; + case 16: sqlite3_result_int64(ctx, c->window_bytes); break; default: sqlite3_result_null(ctx); break; } return SQLITE_OK; diff --git a/test/postgresql/65_payload_window_cap.sql b/test/postgresql/65_payload_window_cap.sql index 02b9822..3b4e781 100644 --- a/test/postgresql/65_payload_window_cap.sql +++ b/test/postgresql/65_payload_window_cap.sql @@ -154,6 +154,63 @@ SELECT d.windows AS w3, d.nrows AS c3, d.last_wm AS wm3, SELECT (:fail::int + 1) AS fail \gset \endif +-- A stateless caller fetches one chunk per call and resumes through resume_*, which is +-- how the /check job pages. Each call builds its state fresh, so the budget spent has to +-- be carried back in through resume_window_bytes or the count restarts every time and +-- the cap never fires at all -- the whole window comes out uncapped however long it runs. +-- The budget here spans several chunks, so it can only be reached by accumulating. +CREATE FUNCTION pg_temp.drain_paged(cap bigint) +RETURNS TABLE (windows int, nrows bigint, last_wm bigint, contiguous boolean) AS $$ +DECLARE + since bigint := 0; wm bigint; w int := 0; nr bigint := 0; ok boolean := true; + rdbv bigint; rseq bigint; rfrag bigint; spent bigint; + capped boolean; final_seen boolean; first_dbv bigint; last_dbv bigint; n int; + r record; +BEGIN + SELECT max(watermark_db_version) INTO wm FROM cloudsync_payload_chunks(0, NULL, NULL, false); + LOOP + rdbv := NULL; rseq := NULL; rfrag := NULL; spent := 0; + capped := false; final_seen := false; n := 0; + LOOP + IF n = 0 THEN + SELECT * INTO r FROM cloudsync_payload_chunks(since, NULL, NULL, false, NULL, NULL, NULL, cap, 0) LIMIT 1; + ELSE + SELECT * INTO r FROM cloudsync_payload_chunks(NULL, NULL, wm, false, rdbv, rseq, rfrag, cap, spent) LIMIT 1; + END IF; + EXIT WHEN r IS NULL; + nr := nr + r.rows; + rdbv := r.next_db_version; rseq := r.next_seq; rfrag := r.next_frag_offset; + final_seen := r.is_final; capped := r.window_capped; spent := r.window_bytes; + IF n = 0 THEN first_dbv := r.db_version_min; END IF; + last_dbv := r.db_version_max; + n := n + 1; + EXIT WHEN final_seen OR n > 500; + END LOOP; + EXIT WHEN n = 0; + -- Windows abut and advance, and the budget is only ever reached across calls. + IF first_dbv <> since + 1 OR last_dbv <= since OR spent > cap + 3 * 262144 THEN ok := false; END IF; + w := w + 1; + IF capped THEN since := last_dbv; ELSE since := wm; END IF; + EXIT WHEN NOT capped OR w > 100; + END LOOP; + RETURN QUERY SELECT w, nr, since, ok; +END $$ LANGUAGE plpgsql; + +-- Fresh totals: the fragment case above added rows after the first baseline was taken. +SELECT sum(rows) AS all_rows, max(watermark_db_version) AS all_wm + FROM cloudsync_payload_chunks(0, NULL, NULL, false) \gset + +SELECT d.windows AS w4, d.nrows AS c4, d.last_wm AS wm4, + (d.windows > 1 AND d.contiguous AND d.nrows = :all_rows::bigint + AND d.last_wm = :all_wm::bigint) AS cap4_ok + FROM pg_temp.drain_paged(600000) d \gset +\if :cap4_ok +\echo [PASS] (:testid) a budget spanning chunks caps a stream paged one chunk per call (:w4 windows) +\else +\echo [FAIL] (:testid) windows=:w4 rows=:c4/:all_rows wm=:wm4/:all_wm +SELECT (:fail::int + 1) AS fail \gset +\endif + \connect postgres \ir helper_psql_conn_setup.sql DROP DATABASE IF EXISTS cloudsync_test_65; diff --git a/test/review_regressions.c b/test/review_regressions.c index 4ede5d1..d1549a6 100644 --- a/test/review_regressions.c +++ b/test/review_regressions.c @@ -846,6 +846,86 @@ static void test_payload_window_cap(void) { CHECK(close_db(db) == SQLITE_OK); } +// A stateless caller fetches one chunk per call and resumes through resume_*, which is +// how the /check job pages. Each call is a fresh scan, so the budget spent has to be +// carried back in through resume_window_bytes or the count restarts every time and the +// cap never fires at all -- the whole window comes out uncapped however long it runs. +static void test_payload_window_cap_paged(void) { + sqlite3 *db = open_db(); + CHECK(sql(db, "CREATE TABLE t(id TEXT PRIMARY KEY NOT NULL, v BLOB);" + "SELECT cloudsync_init('t');" + "SELECT cloudsync_set('payload_max_chunk_size','262144');") == SQLITE_OK); + for (int i = 0; i < 60; i++) { + char buf[128]; + snprintf(buf, sizeof(buf), "INSERT INTO t VALUES('r%03d', randomblob(20000));", i); + CHECK(sql(db, buf) == SQLITE_OK); + } + + sqlite3_stmt *vm = NULL; + int64_t watermark = 0, total_rows = 0; + CHECK(sqlite3_prepare_v2(db, "SELECT max(CASE WHEN is_final THEN watermark_db_version END), sum(rows) " + "FROM cloudsync_payload_chunks(0,NULL,NULL,false)", -1, &vm, NULL) == SQLITE_OK); + if (sqlite3_step(vm) == SQLITE_ROW) { watermark = sqlite3_column_int64(vm, 0); total_rows = sqlite3_column_int64(vm, 1); } + sqlite3_finalize(vm); + CHECK(watermark == 60 && total_rows > 0); + + // A budget spanning several chunks: it can only be reached by accumulating across + // calls, so this is exactly the case a per-call counter cannot cap. + const int64_t budget = 600000; + int64_t since = 0, seen_rows = 0; + int windows = 0; + while (windows < 100) { + int64_t rdbv = 0, rseq = 0, rfrag = 0, spent = 0; + bool capped = false, final_seen = false; + int chunks = 0; + int64_t window_first = 0, window_last = 0; + while (chunks < 500) { + char q[640]; + if (chunks == 0) + snprintf(q, sizeof(q), + "SELECT rows, next_db_version, next_seq, next_frag_offset, is_final, window_capped, " + "window_bytes, db_version_min, db_version_max " + "FROM cloudsync_payload_chunks(%lld,NULL,NULL,false,NULL,NULL,NULL,%lld,0) LIMIT 1", + (long long)since, (long long)budget); + else + snprintf(q, sizeof(q), + "SELECT rows, next_db_version, next_seq, next_frag_offset, is_final, window_capped, " + "window_bytes, db_version_min, db_version_max " + "FROM cloudsync_payload_chunks(NULL,NULL,%lld,false,%lld,%lld,%lld,%lld,%lld) LIMIT 1", + (long long)watermark, (long long)rdbv, (long long)rseq, (long long)rfrag, + (long long)budget, (long long)spent); + vm = NULL; + CHECK(sqlite3_prepare_v2(db, q, -1, &vm, NULL) == SQLITE_OK); + if (sqlite3_step(vm) != SQLITE_ROW) { sqlite3_finalize(vm); break; } + seen_rows += sqlite3_column_int64(vm, 0); + rdbv = sqlite3_column_int64(vm, 1); + rseq = sqlite3_column_int64(vm, 2); + rfrag = sqlite3_column_int64(vm, 3); + final_seen = sqlite3_column_int(vm, 4) != 0; + capped = sqlite3_column_int(vm, 5) != 0; + spent = sqlite3_column_int64(vm, 6); + if (chunks == 0) window_first = sqlite3_column_int64(vm, 7); + window_last = sqlite3_column_int64(vm, 8); + sqlite3_finalize(vm); + chunks++; + if (final_seen) break; + } + if (chunks == 0) break; + CHECK(window_first == since + 1); // windows abut: no gap, no overlap + CHECK(window_last > since); // and each one advances + // The budget is reached by accumulating across calls, never within one. + CHECK(spent <= budget + 3 * 262144); + since = capped ? window_last : watermark; + windows++; + if (!capped) break; + } + CHECK(windows > 1); // the budget really did split the stream + CHECK(since == watermark); // and the drain reached the end + CHECK(seen_rows == total_rows); // covering every row exactly once + CHECK(close_db(db) == SQLITE_OK); +} + + int main(void) { CHECK(sqlite3_config(SQLITE_CONFIG_GETMALLOC, &memory) == SQLITE_OK); sqlite3_mem_methods faults = memory; @@ -869,6 +949,7 @@ int main(void) { test_block_not_null_payload(); test_block_group_atomicity(); test_payload_window_cap(); + test_payload_window_cap_paged(); test_refill_error(); test_block_oom(); cloudsync_memory_finalize(); From 831b157641fe26ed88b865fae4c062168aaeeca8 Mon Sep 17 00:00:00 2001 From: Andrea Donetti Date: Fri, 25 Sep 2026 17:56:30 -0600 Subject: [PATCH 10/14] fix(postgres): guard the new budget arguments with PG_NARGS cloudsync_payload_chunks read argument slots 7 and 8 unconditionally. Between installing the 1.2 binary and running ALTER EXTENSION cloudsync UPDATE, the 1.1 SQL definition still points at this function and calls it with seven arguments, so fcinfo->args is sized for seven and those reads run off the end of it. That window is reachable on any tenant where the two steps are separated, and today's servers make exactly that seven-argument call. In practice the read lands in allocator padding and comes back as NULL, which is why it does not show up as a crash -- but the value read is whatever follows the array, and a garbage budget would cap windows at nonsense sizes rather than fail visibly. Both reads are now guarded, matching how PG_NARGS() is already used elsewhere in this file. Verified with a seven-parameter SQL declaration bound to the same C symbol, which is the upgrade window reproduced exactly. Also drops a stray ")" from the API.md heading. Co-Authored-By: Claude Opus 5 (1M context) --- API.md | 2 +- src/postgresql/cloudsync_postgresql.c | 8 ++++++-- 2 files changed, 7 insertions(+), 3 deletions(-) diff --git a/API.md b/API.md index 72d8c79..0e90991 100644 --- a/API.md +++ b/API.md @@ -567,7 +567,7 @@ SELECT cloudsync_payload_blob_checked( --- -### `cloudsync_payload_chunks([since_db_version], [filter_site_id], [until_db_version], [exclude_filter_site_id], [resume_db_version], [resume_seq], [resume_frag_offset], [max_window_bytes]), [resume_window_bytes])` +### `cloudsync_payload_chunks([since_db_version], [filter_site_id], [until_db_version], [exclude_filter_site_id], [resume_db_version], [resume_seq], [resume_frag_offset], [max_window_bytes], [resume_window_bytes])` **Description:** Generates sync payloads as a stream of transport-sized chunks. It is the chunk-aware evolution of [`cloudsync_payload_encode()`](#cloudsync_payload_encodetbl-pk-col_name-col_value-col_version-db_version-site_id-cl-seq), designed for large rowsets and for single BLOB/TEXT values that are larger than the configured chunk size. diff --git a/src/postgresql/cloudsync_postgresql.c b/src/postgresql/cloudsync_postgresql.c index b95343d..4e63bc4 100644 --- a/src/postgresql/cloudsync_postgresql.c +++ b/src/postgresql/cloudsync_postgresql.c @@ -1400,13 +1400,17 @@ Datum cloudsync_payload_chunks(PG_FUNCTION_ARGS) { int64 resume_dbv = PG_ARGISNULL(4) ? 0 : PG_GETARG_INT64(4); int64 resume_seq = PG_ARGISNULL(5) ? 0 : PG_GETARG_INT64(5); int64 resume_frag = PG_ARGISNULL(6) ? 0 : PG_GETARG_INT64(6); + // Both budget arguments are guarded by PG_NARGS(): between installing the 1.2 + // binary and running ALTER EXTENSION cloudsync UPDATE, the 1.1 SQL definition + // still points at this function and calls it with seven arguments. Reading the + // eighth and ninth slots then runs off the end of fcinfo->args. // Cap the whole prepared window, not one chunk: <= 0 and NULL both mean no cap. - int64 window_cap = PG_ARGISNULL(7) ? 0 : PG_GETARG_INT64(7); + int64 window_cap = (PG_NARGS() > 7 && !PG_ARGISNULL(7)) ? PG_GETARG_INT64(7) : 0; st->max_window_bytes = (window_cap > 0) ? window_cap : 0; // Budget already spent by earlier calls of this window. State is created fresh // per call, so a caller paging one chunk at a time seeds it here; a caller // draining the window in one call leaves it at 0 and it accumulates. - int64 window_spent = PG_ARGISNULL(8) ? 0 : PG_GETARG_INT64(8); + int64 window_spent = (PG_NARGS() > 8 && !PG_ARGISNULL(8)) ? PG_GETARG_INT64(8) : 0; st->window_bytes = (window_spent > 0) ? window_spent : 0; // Site filter resolution: // exclude=true -> all sites except filter_site_id (CHECK path); site required From 67c1962d98f8b1230a9a32b9eb394d38fb18e445 Mon Sep 17 00:00:00 2001 From: Andrea Donetti Date: Fri, 25 Sep 2026 19:37:12 -0600 Subject: [PATCH 11/14] docs: show the stateless paging loop for cloudsync_payload_chunks The usage examples covered only the four original arguments. The resume_* trio had never been shown at all, and the budget arguments are unusable without the loop: each query is a new scan, so window_bytes has to be echoed back as resume_window_bytes or the budget restarts every time. Co-Authored-By: Claude Opus 5 (1M context) --- API.md | 37 +++++++++++++++++++++++++++++++++++++ 1 file changed, 37 insertions(+) diff --git a/API.md b/API.md index 0e90991..50b330d 100644 --- a/API.md +++ b/API.md @@ -636,8 +636,34 @@ WHERE since_db_version = 100 AND site_id = cloudsync_uuid_blob('0190a1b2-c3d4-7e5f-8a9b-001122334455') AND exclude_filter_site_id = 1 ORDER BY chunk_index; + +-- Stateless paging: one chunk per query, resuming where the previous one stopped, +-- under a window budget. Carry next_* back as resume_*, and window_bytes back as +-- resume_window_bytes -- each query is a new scan, so without the echo the budget +-- restarts every time and a budget larger than one chunk is never reached. +SELECT payload, watermark_db_version, is_final, window_capped, window_bytes, + next_db_version, next_seq, next_frag_offset +FROM cloudsync_payload_chunks +WHERE since_db_version = 100 + AND max_window_bytes = 134217728 +LIMIT 1; + +-- ...then for each following chunk of the same window, with the values the +-- previous query returned: +SELECT payload, watermark_db_version, is_final, window_capped, window_bytes, + next_db_version, next_seq, next_frag_offset +FROM cloudsync_payload_chunks +WHERE until_db_version = 200 -- the watermark the first chunk reported + AND resume_db_version = 142 -- next_db_version + AND resume_seq = 0 -- next_seq + AND resume_frag_offset = 0 -- next_frag_offset + AND max_window_bytes = 134217728 -- unchanged for the whole window + AND resume_window_bytes = 5242880 -- window_bytes +LIMIT 1; ``` +Stop when `is_final` is true. If `window_capped` is also true the budget ended the window early: checkpoint at that chunk's `watermark_db_version` and start a new window from it, with `resume_window_bytes` back at `0`. + **PostgreSQL usage:** `cloudsync_payload_chunks` is exposed as a set-returning function with optional arguments: ```sql @@ -652,6 +678,17 @@ FROM cloudsync_payload_chunks(100, cloudsync_siteid(), 200); -- /check download: all changes EXCEPT the requesting peer's site SELECT * FROM cloudsync_payload_chunks(100, cloudsync_uuid_blob('0190a1b2-c3d4-7e5f-8a9b-001122334455'), NULL, true); + +-- Stateless paging under a window budget: one chunk per call, carrying next_* back as +-- resume_* and window_bytes back as resume_window_bytes. +SELECT * FROM cloudsync_payload_chunks(100, NULL, NULL, false, + NULL, NULL, NULL, -- fresh stream + 134217728, 0) LIMIT 1; + +SELECT * FROM cloudsync_payload_chunks(NULL, NULL, 200, false, + 142, 0, 0, -- next_db_version, next_seq, next_frag_offset + 134217728, 5242880) -- same budget, window_bytes echoed back +LIMIT 1; ``` **Apply example:** From 54580196aafb2f379530fd491fb9c8493b55072f Mon Sep 17 00:00:00 2001 From: Andrea Donetti Date: Fri, 25 Sep 2026 19:47:01 -0600 Subject: [PATCH 12/14] docs: qualify what receive.complete=true actually means API.md told callers to loop while complete is false, which is right, but implied the converse: that true means nothing remains. It does not. A drain that ends on a 202 while the server is still preparing reports complete: true with rows 0 (issue #80), and so would a window the server bounded deliberately. Co-Authored-By: Claude Opus 5 (1M context) --- API.md | 2 ++ 1 file changed, 2 insertions(+) diff --git a/API.md b/API.md index 50b330d..0c395ef 100644 --- a/API.md +++ b/API.md @@ -889,6 +889,8 @@ By default this function **drains all currently-available chunks** in one call. SELECT cloudsync_network_receive_changes(5) ->> '$.receive.complete'; ``` +`receive.complete` is `false` only when the stream is *known* to have more pending. `true` means this drain delivered everything that was ready — it is not a guarantee that the server holds nothing further. It is also reported when the server is still preparing a download and answers with no artifact (see issue #80), and would be reported for a window the server bounded deliberately. So loop while it is `false`, but treat `true` as "nothing more was ready just now" and keep syncing on your normal schedule rather than concluding the database is fully up to date. + The drain position (the per-stream page cursor) is held **in memory** on the network context, so a capped drain resumes where it left off on the next call — the caller does not manage any cursor; it just loops while `receive.complete` is `false`. If the connection is closed or the process restarts mid-drain, or a call fails, the next call safely restarts the drain from the beginning of the stream: already-applied chunks are re-downloaded and re-applied idempotently, so **no rows are skipped** — only redundant download is incurred. This is safe because the durable receive checkpoint (`check_dbversion`/`check_seq`) stays fixed for the whole stream and only advances once the stream has been **fully** applied. A stream whose final chunk arrives while a fragmented value it delivered is still incomplete fails with an error instead of advancing. If the network is misconfigured or the remote server is unreachable, the function raises a SQL error. If the received payload cannot be applied locally (for example because of an unknown schema hash), the error is returned as a `receive.error` field in the JSON response. If the server reports an unresolved failed check job (e.g. an `encode_changes` failure), that failure is forwarded as a `receive.lastFailure` object. From 1403d3f4a0ec6873aafbd06af913f08950e6ac95 Mon Sep 17 00:00:00 2001 From: Andrea Donetti Date: Mon, 28 Sep 2026 10:34:14 -0600 Subject: [PATCH 13/14] docs: release 1.2.0 Marks the unreleased section as 1.2.0, dated today, and corrects the extension version entry: cloudsync_payload_chunks() gained two arguments and two output columns, not one of each, after the budget had to be carried across calls. Co-Authored-By: Claude Opus 5 (1M context) --- CHANGELOG.md | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index ee59546..ebca23c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,7 +4,7 @@ All notable changes to this project will be documented in this file. The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/). -## [Unreleased] +## [1.2.0] - 2026-09-28 ### Added @@ -14,7 +14,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/). ### Changed - **SQLite: an explicit `NULL` for `resume_db_version` on `cloudsync_payload_chunks()` now means "not given"**, matching what it has always meant for `filter_site_id` and on PostgreSQL. It was previously read as a resume point of database version 0, which silently ignored `since_db_version` and restarted the scan at the beginning of the window. Reaching a later argument requires passing `NULL` for the ones before it, so this is easy to hit: `cloudsync_payload_chunks(100, NULL, NULL, false, NULL, NULL, NULL, 1048576)` used to replay from the start of the history instead of resuming after version 100. -- **The PostgreSQL extension version moves to `1.2`.** `cloudsync_payload_chunks()` gained an argument and an output column, so existing deployments need `ALTER EXTENSION cloudsync UPDATE;` after installing the new binary. The upgrade script replaces the function: a `CREATE OR REPLACE` cannot change a return type, and leaving the old seven-argument version in place would make every existing call ambiguous against the new eight-argument one. +- **The PostgreSQL extension version moves to `1.2`.** `cloudsync_payload_chunks()` gained two arguments and two output columns, so existing deployments need `ALTER EXTENSION cloudsync UPDATE;` after installing the new binary. The upgrade script replaces the function: a `CREATE OR REPLACE` cannot change a return type, and leaving the old seven-argument version in place would make every existing call ambiguous against the new eight-argument one. ### Fixed From ea8b5ea0757e2043c343fd9a65aa074366c78cb4 Mon Sep 17 00:00:00 2001 From: Andrea Donetti Date: Mon, 28 Sep 2026 11:11:55 -0600 Subject: [PATCH 14/14] ci: pull the Supabase base image from Docker Hub The supabase docker jobs failed twice on "toomanyrequests: Data limit exceeded" while resolving FROM public.ecr.aws/supabase/postgres. Anonymous ECR Public pulls are charged to a quota keyed on the source IP, which on GitHub-hosted runners is a NAT pool shared with every other Actions user. Supabase publishes the same image to Docker Hub. Both tags resolve to the same manifest index digest -- 17.6.1.170 to sha256:6087b151...6728f and 15.14.1.170 to sha256:bc165447...70b02 -- with linux/amd64 and linux/arm64 on each, so the published images are unchanged. The docker-publish job already logs in to Docker Hub before it builds, so these pulls are now authenticated and counted against our own account. make postgres-supabase-build is unaffected: its sed already rewrites both spellings of the FROM line to the running Supabase CLI image tag, which keeps its public.ecr.aws name so the CLI still picks up the rebuilt image. Co-Authored-By: Claude Opus 5 (1M context) --- docker/postgresql/Dockerfile.supabase | 6 ++++-- docker/postgresql/Dockerfile.supabase.release | 6 +++++- 2 files changed, 9 insertions(+), 3 deletions(-) diff --git a/docker/postgresql/Dockerfile.supabase b/docker/postgresql/Dockerfile.supabase index 1010702..07360fb 100644 --- a/docker/postgresql/Dockerfile.supabase +++ b/docker/postgresql/Dockerfile.supabase @@ -46,8 +46,10 @@ RUN mkdir -p /tmp/cloudsync-artifacts/lib /tmp/cloudsync-artifacts/extension && cp /tmp/cloudsync/src/postgresql/migrations/cloudsync--*--*.sql /tmp/cloudsync-artifacts/extension/; \ fi -# Runtime image based on Supabase Postgres -FROM public.ecr.aws/supabase/postgres:${SUPABASE_POSTGRES_TAG} +# Runtime image based on Supabase Postgres. Docker Hub and public.ecr.aws carry +# the same image; see Dockerfile.supabase.release for why this one is used. +# make postgres-supabase-build rewrites this line to the running CLI image tag. +FROM supabase/postgres:${SUPABASE_POSTGRES_TAG} # Extension version (derived from src/cloudsync.h by the Makefile and passed in # as a build arg); used only for the image label. diff --git a/docker/postgresql/Dockerfile.supabase.release b/docker/postgresql/Dockerfile.supabase.release index 22fde84..1c98012 100644 --- a/docker/postgresql/Dockerfile.supabase.release +++ b/docker/postgresql/Dockerfile.supabase.release @@ -19,7 +19,11 @@ # ARG SUPABASE_POSTGRES_TAG=17.6.1.170 -FROM public.ecr.aws/supabase/postgres:${SUPABASE_POSTGRES_TAG} +# Supabase publishes the same image to Docker Hub and to public.ecr.aws; both +# tags resolve to identical manifest digests. Docker Hub is used here because CI +# already authenticates to it, while anonymous ECR Public pulls are charged to a +# data quota shared by every GitHub Actions runner. +FROM supabase/postgres:${SUPABASE_POSTGRES_TAG} ARG CLOUDSYNC_VERSION ARG TARGETARCH