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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
66 changes: 58 additions & 8 deletions API.md

Large diffs are not rendered by default.

6 changes: 6 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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. 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

- **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

Expand Down
2 changes: 1 addition & 1 deletion src/cloudsync.h
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
8 changes: 6 additions & 2 deletions src/postgresql/cloudsync.sql.in
Original file line number Diff line number Diff line change
Expand Up @@ -156,7 +156,9 @@ 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,
resume_window_bytes bigint DEFAULT NULL
)
RETURNS TABLE (
payload bytea,
Expand All @@ -169,7 +171,9 @@ RETURNS TABLE (
next_db_version bigint,
next_seq bigint,
next_frag_offset bigint,
is_final boolean
is_final boolean,
window_capped boolean,
window_bytes bigint
)
AS 'MODULE_PATHNAME', 'cloudsync_payload_chunks'
LANGUAGE C VOLATILE;
Expand Down
62 changes: 59 additions & 3 deletions src/postgresql/cloudsync_postgresql.c
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -1228,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);
Expand Down Expand Up @@ -1317,6 +1327,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};
Expand All @@ -1334,6 +1351,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;
Expand Down Expand Up @@ -1380,6 +1400,18 @@ 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_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_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
// filter given -> only that site
Expand Down Expand Up @@ -1472,7 +1504,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;
Expand All @@ -1497,8 +1531,28 @@ 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) {
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[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));
Expand All @@ -1510,6 +1564,8 @@ 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);
outvals[12] = Int64GetDatum(st->window_bytes);
HeapTuple outtup = heap_form_tuple(st->outdesc, outvals, outnulls);
SRF_RETURN_NEXT(funcctx, HeapTupleGetDatum(outtup));
}
Expand Down
50 changes: 50 additions & 0 deletions src/postgresql/migrations/cloudsync--1.1--1.2.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,50 @@
-- CloudSync PostgreSQL extension upgrade: 1.1 -> 1.2
--
-- Bounds what one call to cloudsync_payload_chunks() prepares:
-- * 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, 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.
--
-- 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,
resume_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,
window_bytes bigint
)
AS 'MODULE_PATHNAME', 'cloudsync_payload_chunks'
LANGUAGE C VOLATILE;
Loading
Loading