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
41 changes: 41 additions & 0 deletions packages/zero-cache/src/config/normalize.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ function configWith(litestream: Partial<ZeroConfig['litestream']>): ZeroConfig {
changeStreamer: {
port: 4849,
address: 'localhost',
pgChangeLogEnabled: true,
sqliteChangeLogMode: 'off',
sqliteChangeLogReadPercent: 0,
sqliteChangeLogColdReadPercent: 0,
Expand Down Expand Up @@ -128,6 +129,46 @@ describe('config/normalize litestream v5 gating', () => {
});

describe('config/normalize SQLite change log', () => {
test('PG change log is enabled by default configuration', () => {
const config = configWith({});

expect(config.changeStreamer.pgChangeLogEnabled).toBe(true);
expect(() => assertNormalized(config)).not.toThrow();
});

test('disabling the PG change log requires authoritative SQLite and v5 backup settings', () => {
const config = configWith({});
config.changeStreamer.pgChangeLogEnabled = false;

expect(() => assertNormalized(config)).toThrow(
'requires --change-streamer-sqlite-change-log-mode=serve',
);

config.changeStreamer.sqliteChangeLogMode = 'serve';
expect(() => assertNormalized(config)).toThrow(
'requires --change-streamer-sqlite-change-log-read-percent=100',
);

config.changeStreamer.sqliteChangeLogReadPercent = 100;
expect(() => assertNormalized(config)).toThrow(
'requires --change-streamer-sqlite-change-log-cold-read-percent=100',
);

config.changeStreamer.sqliteChangeLogColdReadPercent = 100;
expect(() => assertNormalized(config)).toThrow(
'requires a litestream v5 backup',
);

Object.assign(config.litestream, {
backupURL: 's3://bucket/replica',
backupUsingV5: true,
restoreUsingV5: true,
executableV5: '/bin/litestream-v5',
vfsQueryExecutable: '/bin/vfs-query',
});
expect(() => assertNormalized(config)).not.toThrow();
});

test('read percentage is only allowed in serve mode', () => {
const config = configWith({});
config.changeStreamer.sqliteChangeLogMode = 'compare';
Expand Down
19 changes: 19 additions & 0 deletions packages/zero-cache/src/config/normalize.ts
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,7 @@ export function assertNormalized(
assert(config.changeStreamer.port, 'missing --change-streamer-port');
assert(config.changeStreamer.address, 'missing --change-streamer-address');
const {
pgChangeLogEnabled,
sqliteChangeLogMode,
sqliteChangeLogReadPercent,
sqliteChangeLogColdReadPercent,
Expand Down Expand Up @@ -97,6 +98,24 @@ export function assertNormalized(
sqliteChangeLogReadPercent > 0 || sqliteChangeLogColdReadPercent === 0,
'--change-streamer-sqlite-change-log-cold-read-percent must be 0 when --change-streamer-sqlite-change-log-read-percent is 0',
);
if (!pgChangeLogEnabled) {
assert(
sqliteChangeLogMode === 'serve',
'--change-streamer-pg-change-log-enabled=false requires --change-streamer-sqlite-change-log-mode=serve',
);
assert(
sqliteChangeLogReadPercent === 100,
'--change-streamer-pg-change-log-enabled=false requires --change-streamer-sqlite-change-log-read-percent=100',
);
assert(
sqliteChangeLogColdReadPercent === 100,
'--change-streamer-pg-change-log-enabled=false requires --change-streamer-sqlite-change-log-cold-read-percent=100',
);
assert(
config.litestream.backupURL && config.litestream.backupUsingV5,
'--change-streamer-pg-change-log-enabled=false requires a litestream v5 backup',
);
}
for (const [flag, value] of [
['retention-ms', sqliteChangeLogRetentionMs],
['read-batch-rows', sqliteChangeLogReadBatchRows],
Expand Down
18 changes: 18 additions & 0 deletions packages/zero-cache/src/config/zero-config.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -905,6 +905,24 @@ test('--enable-query-covering can be disabled', () => {
expect(config.enableQueryCovering).toBe(false);
});

test('PG change log is enabled by default and can be disabled by env', () => {
const defaults = parseOptionsAdvanced(zeroOptions, {
envNamePrefix: 'ZERO_',
allowUnknown: false,
allowPartial: true,
env: {},
}).config;
const disabled = parseOptionsAdvanced(zeroOptions, {
envNamePrefix: 'ZERO_',
allowUnknown: false,
allowPartial: true,
env: {ZERO_CHANGE_STREAMER_PG_CHANGE_LOG_ENABLED: 'false'},
}).config;

expect(defaults.changeStreamer.pgChangeLogEnabled).toBe(true);
expect(disabled.changeStreamer.pgChangeLogEnabled).toBe(false);
});

test('legacy queries are disabled by default', () => {
const {config} = parseOptionsAdvanced(zeroOptions, {
envNamePrefix: 'ZERO_',
Expand Down
10 changes: 10 additions & 0 deletions packages/zero-cache/src/config/zero-config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -737,6 +737,16 @@ export const zeroOptions = {
hidden: true,
},

pgChangeLogEnabled: {
type: v.boolean().default(true),
desc: [
`Whether the legacy Postgres change log remains authoritative for`,
`stream initialization, persistence, catchup, and upstream ACKs.`,
`Disabling it requires SQLite serve mode at 100 percent and a v5 backup.`,
],
hidden: true,
},

sqliteChangeLogReadPercent: {
type: v.number().default(0),
desc: [
Expand Down
21 changes: 12 additions & 9 deletions packages/zero-cache/src/server/change-streamer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,7 @@ export default async function runWorker(
backPressureLimitHeapProportion,
flowControlConsensusTimeoutProportion,
flowControlSlowSubscriberGracePeriodSeconds,
pgChangeLogEnabled,
sqliteChangeLogMode,
sqliteChangeLogReadPercent,
sqliteChangeLogColdReadPercent,
Expand Down Expand Up @@ -123,7 +124,7 @@ export default async function runWorker(
// purges. This ensures that (this) change-streamer will be able to resume
// from the backup.
let purgeLock =
litestream.backupURL && litestream.executable
pgChangeLogEnabled && litestream.backupURL && litestream.executable
? await new PurgeLocker(lc, shard, changeDB).acquire()
: null;
const restoreOptions = {litestream, constraints: purgeLock ?? undefined};
Expand Down Expand Up @@ -220,6 +221,7 @@ export default async function runWorker(
purgeLock,
autoReset ?? false,
{
pgChangeLogEnabled,
backPressureLimitHeapProportion,
flowControlConsensusTimeoutProportion,
flowControlSlowSubscriberGracePeriodMs:
Expand Down Expand Up @@ -258,14 +260,15 @@ export default async function runWorker(
}
: undefined,
// Compare mode runs both advisory checks. Postgres remains authoritative.
sqliteChangeLogCompare: sqliteChangeLogComparing
? {
replicaFile: replica.file,
comparePercent: sqliteChangeLogComparePercent,
retentionMs: sqliteChangeLogRetentionMs,
readBatchRows: sqliteChangeLogReadBatchRows,
}
: undefined,
sqliteChangeLogCompare:
pgChangeLogEnabled && sqliteChangeLogComparing
? {
replicaFile: replica.file,
comparePercent: sqliteChangeLogComparePercent,
retentionMs: sqliteChangeLogRetentionMs,
readBatchRows: sqliteChangeLogReadBatchRows,
}
: undefined,
// Slice 11 lands dark by default: serve mode constructs the stable
// router, while readPercent=0 keeps every catchup on PG and emits
// eligibility metrics before any canary traffic is enabled.
Expand Down
Loading
Loading