diff --git a/API.md b/API.md index 6322be78..0c395ef3 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-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). 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). +**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), 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])` +### `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. @@ -589,6 +589,13 @@ 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. + + 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: @@ -600,7 +607,11 @@ 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`, **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`. | +| `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: @@ -625,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 @@ -641,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:** @@ -660,7 +708,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-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. @@ -779,7 +827,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-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. @@ -841,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. diff --git a/CHANGELOG.md b/CHANGELOG.md index 7e1b43ae..ebca23ce 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,11 +4,17 @@ 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 - **`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 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 diff --git a/docker/postgresql/Dockerfile.supabase b/docker/postgresql/Dockerfile.supabase index 10107022..07360fba 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 22fde841..1c980125 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 diff --git a/src/cloudsync.h b/src/cloudsync.h index ffafc4b6..a1476056 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 92bb6f42..9df66b71 100644 --- a/src/postgresql/cloudsync.sql.in +++ b/src/postgresql/cloudsync.sql.in @@ -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, @@ -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; diff --git a/src/postgresql/cloudsync_postgresql.c b/src/postgresql/cloudsync_postgresql.c index 9f7145a6..4e63bc49 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; @@ -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); @@ -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}; @@ -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; @@ -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 @@ -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; @@ -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)); @@ -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)); } 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 00000000..fa280b02 --- /dev/null +++ b/src/postgresql/migrations/cloudsync--1.1--1.2.sql @@ -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; diff --git a/src/sqlite/cloudsync_sqlite.c b/src/sqlite/cloudsync_sqlite.c index 8345dd81..3e6e236e 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,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, " - "resume_db_version HIDDEN, resume_seq HIDDEN, resume_frag_offset HIDDEN)"); + // 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, resume_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 +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(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(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) { @@ -1219,6 +1237,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); @@ -1261,6 +1283,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); @@ -1273,6 +1302,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++; @@ -1305,9 +1337,37 @@ 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; + 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 +1413,29 @@ 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 + } + // 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 @@ -1445,7 +1525,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 +1550,8 @@ 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; + 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 new file mode 100644 index 00000000..3b4e7814 --- /dev/null +++ b/test/postgresql/65_payload_window_cap.sql @@ -0,0 +1,216 @@ +-- 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 +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 + +-- 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; + +-- 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 sum(rows) AS base_rows, sum(payload_size) AS base_bytes, + max(watermark_db_version) AS base_wm, + (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 rows=:base_rows 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_from(start_since bigint, cap bigint) +RETURNS TABLE (windows int, nrows bigint, maxbytes bigint, last_wm bigint, contiguous boolean) AS $$ +DECLARE + since bigint := start_since; + w int := 0; + c bigint := 0; + b bigint := 0; + ok boolean := true; + r record; +BEGIN + LOOP + 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; + -- 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.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_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 +\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.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_from(0, 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 rows=:c2/:base_rows maxwindow=:b2 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 + +-- 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 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 + +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 + +-- 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/postgresql/full_test.sql b/test/postgresql/full_test.sql index fe21bd0f..6baafc10 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' diff --git a/test/review_regressions.c b/test/review_regressions.c index 14e27fd1..d1549a69 100644 --- a/test/review_regressions.c +++ b/test/review_regressions.c @@ -690,6 +690,242 @@ 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); + } + // 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_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_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_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_rows = 0; + int windows = 0; + bool capped = true; + while (capped && windows < 100) { + char q[512]; + snprintf(q, sizeof(q), + "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 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_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 + // 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); + + // 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); +} + +// 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; @@ -712,6 +948,8 @@ int main(void) { test_block_migration_orphan(); 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();