Skip to content

Commit 35c0f50

Browse files
authored
feat: bound what one cloudsync_payload_chunks() call prepares (#77)
A single cloudsync_payload_chunks() call drained the whole window, so a large backlog produced one unbounded run of work -- the shape behind the prepare stall, where a job kept reaching its deadline before it finished. Two optional arguments now bound it, declared last so the existing positional arguments 1..7 keep their meaning: * max_window_bytes caps the uncompressed bytes one call prepares * resume_window_bytes carries the budget spent so far, for a stream fetched one chunk per call and two outputs report the outcome: window_capped, true when the scan stopped on the budget rather than because the window was drained, and window_bytes, the budget spent. Without max_window_bytes the function behaves exactly as before. A capped call always stops on a db_version boundary, so the receive checkpoint stays on a complete version, and fragment payloads count toward the budget. Continuation is the caller's job: nothing loops automatically, and the API docs qualify what receive.complete=true actually means (#83). Ships as 1.2.0, with the PostgreSQL migration cloudsync--1.1--1.2.sql. The migration drops the old 7-argument function before creating the new one: CREATE OR REPLACE cannot change a return type, and a defaulted new parameter would leave every 7-argument call ambiguous. Also pulls the Supabase base image from Docker Hub rather than public.ecr.aws. Both tags resolve to the same manifest digest, and CI already authenticates to Docker Hub, so the supabase image jobs stop failing on the shared anonymous ECR Public data quota. 🤖 Generated with [Claude Code](https://claude.com/claude-code)
1 parent 9c56e9a commit 35c0f50

12 files changed

Lines changed: 738 additions & 25 deletions

‎API.md‎

Lines changed: 58 additions & 8 deletions
Large diffs are not rendered by default.

‎CHANGELOG.md‎

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,11 +4,17 @@ All notable changes to this project will be documented in this file.
44

55
The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/).
66

7-
## [Unreleased]
7+
## [1.2.0] - 2026-09-28
88

99
### Added
1010

1111
- **`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.
12+
- **`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.
13+
14+
### Changed
15+
16+
- **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.
17+
- **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.
1218

1319
### Fixed
1420

‎docker/postgresql/Dockerfile.supabase‎

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -46,8 +46,10 @@ RUN mkdir -p /tmp/cloudsync-artifacts/lib /tmp/cloudsync-artifacts/extension &&
4646
cp /tmp/cloudsync/src/postgresql/migrations/cloudsync--*--*.sql /tmp/cloudsync-artifacts/extension/; \
4747
fi
4848

49-
# Runtime image based on Supabase Postgres
50-
FROM public.ecr.aws/supabase/postgres:${SUPABASE_POSTGRES_TAG}
49+
# Runtime image based on Supabase Postgres. Docker Hub and public.ecr.aws carry
50+
# the same image; see Dockerfile.supabase.release for why this one is used.
51+
# make postgres-supabase-build rewrites this line to the running CLI image tag.
52+
FROM supabase/postgres:${SUPABASE_POSTGRES_TAG}
5153

5254
# Extension version (derived from src/cloudsync.h by the Makefile and passed in
5355
# as a build arg); used only for the image label.

‎docker/postgresql/Dockerfile.supabase.release‎

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,11 @@
1919
#
2020

2121
ARG SUPABASE_POSTGRES_TAG=17.6.1.170
22-
FROM public.ecr.aws/supabase/postgres:${SUPABASE_POSTGRES_TAG}
22+
# Supabase publishes the same image to Docker Hub and to public.ecr.aws; both
23+
# tags resolve to identical manifest digests. Docker Hub is used here because CI
24+
# already authenticates to it, while anonymous ECR Public pulls are charged to a
25+
# data quota shared by every GitHub Actions runner.
26+
FROM supabase/postgres:${SUPABASE_POSTGRES_TAG}
2327

2428
ARG CLOUDSYNC_VERSION
2529
ARG TARGETARCH

‎src/cloudsync.h‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,7 @@
1818
extern "C" {
1919
#endif
2020

21-
#define CLOUDSYNC_VERSION "1.1.4"
21+
#define CLOUDSYNC_VERSION "1.2.0"
2222
// LZ4's block format cannot expand input by more than 255:1, so a compressed payload
2323
// declaring a larger expansion is forged or corrupt (see cloudsync_payload_apply).
2424
#define CLOUDSYNC_PAYLOAD_LZ4_MAX_RATIO 255

‎src/postgresql/cloudsync.sql.in‎

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -156,7 +156,9 @@ CREATE OR REPLACE FUNCTION cloudsync_payload_chunks(
156156
exclude_filter_site_id boolean DEFAULT false,
157157
resume_db_version bigint DEFAULT NULL,
158158
resume_seq bigint DEFAULT NULL,
159-
resume_frag_offset bigint DEFAULT NULL
159+
resume_frag_offset bigint DEFAULT NULL,
160+
max_window_bytes bigint DEFAULT NULL,
161+
resume_window_bytes bigint DEFAULT NULL
160162
)
161163
RETURNS TABLE (
162164
payload bytea,
@@ -169,7 +171,9 @@ RETURNS TABLE (
169171
next_db_version bigint,
170172
next_seq bigint,
171173
next_frag_offset bigint,
172-
is_final boolean
174+
is_final boolean,
175+
window_capped boolean,
176+
window_bytes bigint
173177
)
174178
AS 'MODULE_PATHNAME', 'cloudsync_payload_chunks'
175179
LANGUAGE C VOLATILE;

‎src/postgresql/cloudsync_postgresql.c‎

Lines changed: 59 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1083,6 +1083,12 @@ typedef struct {
10831083
bool eof;
10841084
int64 chunk_index;
10851085
int64 watermark;
1086+
// Window cap (max_window_bytes): payload bytes emitted so far by this scan, and
1087+
// whether it ended because the budget ran out rather than because the window was
1088+
// drained. window_bytes counts across chunks, not within one.
1089+
int64 max_window_bytes;
1090+
int64 window_bytes;
1091+
bool window_capped;
10861092
int max_size;
10871093
int frag_target;
10881094

@@ -1228,6 +1234,10 @@ static bytea *payload_chunks_emit_pg_fragment(PayloadChunksState *st, cloudsync_
12281234
if (rc != DBRES_OK) ereport(ERROR, (errcode(cloudsync_error_sqlstate(data)), errmsg("%s", cloudsync_errmsg(data))));
12291235
rc = cloudsync_payload_encode_final(payload, data);
12301236
if (rc != DBRES_OK) ereport(ERROR, (errcode(cloudsync_error_sqlstate(data)), errmsg("%s", cloudsync_errmsg(data))));
1237+
// A fragment chunk spends the window budget like any other. Fragments are emitted
1238+
// here rather than by the ordinary builder, so without this a history made of
1239+
// oversized values never spends the budget and the cap never fires.
1240+
st->window_bytes += (int64)cloudsync_payload_context_bused(payload);
12311241
int64 blob_size = 0;
12321242
char *blob = cloudsync_payload_blob(payload, &blob_size, rows);
12331243
bytea *result = (bytea *)palloc(VARHDRSZ + blob_size);
@@ -1317,6 +1327,13 @@ static bytea *payload_chunks_build_pg_next(PayloadChunksState *st, cloudsync_con
13171327
}
13181328

13191329
if (cloudsync_payload_context_nrows(payload) > 0 && cloudsync_payload_context_bused(payload) + row_size > (size_t)st->max_size) break;
1330+
// Once the budget is spent, end the chunk at the first db_version boundary so
1331+
// the window can be capped there. Waiting for a chunk to happen to end on a
1332+
// boundary is not enough: when a transaction's rows and a chunk's capacity stay
1333+
// out of step, every chunk ends mid-version and the cap never fires at all.
1334+
if (st->max_window_bytes > 0 && cloudsync_payload_context_nrows(payload) > 0 &&
1335+
st->db_version != *dbv_max &&
1336+
st->window_bytes + (int64)cloudsync_payload_context_bused(payload) >= st->max_window_bytes) break;
13201337

13211338
pgvalue_t *vals[9] = {0};
13221339
text *owned_texts[2] = {0};
@@ -1334,6 +1351,9 @@ static bytea *payload_chunks_build_pg_next(PayloadChunksState *st, cloudsync_con
13341351
cloudsync_memory_free(payload);
13351352
return NULL;
13361353
}
1354+
// Measure the window in the unit max_size is expressed in -- encoded bytes before
1355+
// compression -- so a budget and a chunk size mean the same thing to a caller.
1356+
st->window_bytes += (int64)cloudsync_payload_context_bused(payload);
13371357
int rc = cloudsync_payload_encode_final(payload, data);
13381358
if (rc != DBRES_OK) ereport(ERROR, (errcode(cloudsync_error_sqlstate(data)), errmsg("%s", cloudsync_errmsg(data))));
13391359
int64 blob_size = 0;
@@ -1380,6 +1400,18 @@ Datum cloudsync_payload_chunks(PG_FUNCTION_ARGS) {
13801400
int64 resume_dbv = PG_ARGISNULL(4) ? 0 : PG_GETARG_INT64(4);
13811401
int64 resume_seq = PG_ARGISNULL(5) ? 0 : PG_GETARG_INT64(5);
13821402
int64 resume_frag = PG_ARGISNULL(6) ? 0 : PG_GETARG_INT64(6);
1403+
// Both budget arguments are guarded by PG_NARGS(): between installing the 1.2
1404+
// binary and running ALTER EXTENSION cloudsync UPDATE, the 1.1 SQL definition
1405+
// still points at this function and calls it with seven arguments. Reading the
1406+
// eighth and ninth slots then runs off the end of fcinfo->args.
1407+
// Cap the whole prepared window, not one chunk: <= 0 and NULL both mean no cap.
1408+
int64 window_cap = (PG_NARGS() > 7 && !PG_ARGISNULL(7)) ? PG_GETARG_INT64(7) : 0;
1409+
st->max_window_bytes = (window_cap > 0) ? window_cap : 0;
1410+
// Budget already spent by earlier calls of this window. State is created fresh
1411+
// per call, so a caller paging one chunk at a time seeds it here; a caller
1412+
// draining the window in one call leaves it at 0 and it accumulates.
1413+
int64 window_spent = (PG_NARGS() > 8 && !PG_ARGISNULL(8)) ? PG_GETARG_INT64(8) : 0;
1414+
st->window_bytes = (window_spent > 0) ? window_spent : 0;
13831415
// Site filter resolution:
13841416
// exclude=true -> all sites except filter_site_id (CHECK path); site required
13851417
// filter given -> only that site
@@ -1472,7 +1504,9 @@ Datum cloudsync_payload_chunks(PG_FUNCTION_ARGS) {
14721504
PayloadChunksState *st = (PayloadChunksState *)funcctx->user_fctx;
14731505

14741506
int64 rows = 0, dbv_min = 0, dbv_max = 0;
1475-
bytea *payload = payload_chunks_build_pg_next(st, data, &rows, &dbv_min, &dbv_max);
1507+
// A capped window ends the scan: the rows past it belong to the next window, and
1508+
// emitting them here would contradict the reduced watermark already reported.
1509+
bytea *payload = st->window_capped ? NULL : payload_chunks_build_pg_next(st, data, &rows, &dbv_min, &dbv_max);
14761510
if (!payload) {
14771511
if (st->portal) SPI_cursor_close(st->portal);
14781512
st->portal = NULL;
@@ -1497,8 +1531,28 @@ Datum cloudsync_payload_chunks(PG_FUNCTION_ARGS) {
14971531
next_dbv = st->watermark; next_seq = 0; next_frag = 0; is_final = true;
14981532
}
14991533

1500-
Datum outvals[11];
1501-
bool outnulls[11] = {false,false,false,false,false,false,false,false,false,false,false};
1534+
// Stop once max_window_bytes is spent, at the first db_version boundary, and report
1535+
// the window as an ordinary complete stream ending at that reduced watermark. The
1536+
// caller checkpoints there and asks again, so preparation is bounded with no new
1537+
// resumable state.
1538+
//
1539+
// Two conditions are load-bearing. Never cap mid-value (frag_active) or
1540+
// mid-db_version (next_dbv == dbv_max): the receive cursor must land on a complete
1541+
// db_version or the next request skips the unapplied remainder, since it resumes
1542+
// with db_version > since and no seq. And because the boundary test is what stops
1543+
// us, a db_version larger than the whole budget is still emitted in full --
1544+
// otherwise a window could come out empty and the drain would never advance.
1545+
if (st->max_window_bytes > 0) {
1546+
if (!is_final && !st->frag_active && st->window_bytes >= st->max_window_bytes && next_dbv != dbv_max) {
1547+
st->window_capped = true;
1548+
is_final = true;
1549+
st->watermark = dbv_max;
1550+
next_dbv = st->watermark; next_seq = 0; next_frag = 0;
1551+
}
1552+
}
1553+
1554+
Datum outvals[13];
1555+
bool outnulls[13] = {false,false,false,false,false,false,false,false,false,false,false,false,false};
15021556
outvals[0] = PointerGetDatum(payload);
15031557
outvals[1] = Int64GetDatum(st->chunk_index++);
15041558
outvals[2] = Int64GetDatum(VARSIZE_ANY_EXHDR(payload));
@@ -1510,6 +1564,8 @@ Datum cloudsync_payload_chunks(PG_FUNCTION_ARGS) {
15101564
outvals[8] = Int64GetDatum(next_seq);
15111565
outvals[9] = Int64GetDatum(next_frag);
15121566
outvals[10] = BoolGetDatum(is_final);
1567+
outvals[11] = BoolGetDatum(st->window_capped);
1568+
outvals[12] = Int64GetDatum(st->window_bytes);
15131569
HeapTuple outtup = heap_form_tuple(st->outdesc, outvals, outnulls);
15141570
SRF_RETURN_NEXT(funcctx, HeapTupleGetDatum(outtup));
15151571
}
Lines changed: 50 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,50 @@
1+
-- CloudSync PostgreSQL extension upgrade: 1.1 -> 1.2
2+
--
3+
-- Bounds what one call to cloudsync_payload_chunks() prepares:
4+
-- * new max_window_bytes and resume_window_bytes inputs, declared last so the
5+
-- existing positional arguments 1..7 keep their meaning
6+
-- * new window_capped output, true when the scan stopped on the budget
7+
-- rather than because the window was drained, and window_bytes, the budget
8+
-- spent so far -- pass it back as resume_window_bytes to carry the budget
9+
-- across a stream fetched one chunk per call
10+
--
11+
-- Both are optional: without max_window_bytes the function behaves exactly as
12+
-- it did in 1.1.
13+
--
14+
-- Run automatically by: ALTER EXTENSION cloudsync UPDATE;
15+
16+
-- The old function has to go before the new one is created. CREATE OR REPLACE
17+
-- cannot change a return type ("cannot change return type of existing
18+
-- function"), and because the new parameter has a default, keeping both would
19+
-- leave a 7-argument call matching two candidates -- an ambiguous function
20+
-- call error at every existing call site.
21+
DROP FUNCTION IF EXISTS cloudsync_payload_chunks(bigint, bytea, bigint, boolean, bigint, bigint, bigint);
22+
23+
CREATE OR REPLACE FUNCTION cloudsync_payload_chunks(
24+
since_db_version bigint DEFAULT NULL,
25+
filter_site_id bytea DEFAULT NULL,
26+
until_db_version bigint DEFAULT NULL,
27+
exclude_filter_site_id boolean DEFAULT false,
28+
resume_db_version bigint DEFAULT NULL,
29+
resume_seq bigint DEFAULT NULL,
30+
resume_frag_offset bigint DEFAULT NULL,
31+
max_window_bytes bigint DEFAULT NULL,
32+
resume_window_bytes bigint DEFAULT NULL
33+
)
34+
RETURNS TABLE (
35+
payload bytea,
36+
chunk_index bigint,
37+
payload_size bigint,
38+
rows bigint,
39+
db_version_min bigint,
40+
db_version_max bigint,
41+
watermark_db_version bigint,
42+
next_db_version bigint,
43+
next_seq bigint,
44+
next_frag_offset bigint,
45+
is_final boolean,
46+
window_capped boolean,
47+
window_bytes bigint
48+
)
49+
AS 'MODULE_PATHNAME', 'cloudsync_payload_chunks'
50+
LANGUAGE C VOLATILE;

0 commit comments

Comments
 (0)