From 9876137d4dcaf8f13138eeabc5feb3f3eeb1cacc Mon Sep 17 00:00:00 2001
From: massimoalbarello
Date: Fri, 2 Oct 2026 18:26:11 +0100
Subject: [PATCH 1/5] refactor: centralize sync polling in cron
---
README.md | 6 +-
.../backend/src/services/dashboard/service.ts | 2 -
.../src/routes/_workspace/syncs.$id.tsx | 15 --
.../src/routes/_workspace/syncs.index.tsx | 6 +-
packages/sync/src/db/client.ts | 13 +-
packages/sync/src/db/schema.sql | 2 +-
packages/sync/src/execution/worker.ts | 13 +-
packages/sync/src/http/controller.ts | 1 -
packages/sync/src/models/definition.ts | 2 +-
packages/sync/src/models/sync.ts | 10 +-
.../src/repositories/acquisition/contract.ts | 1 +
.../src/repositories/acquisition/lease.ts | 12 +-
.../src/repositories/acquisition/sqlite.ts | 11 +-
.../sync/src/repositories/assets/sqlite.ts | 4 +-
.../sync/src/repositories/catalog/sqlite.ts | 27 ++--
.../sync/src/repositories/delivery/sqlite.ts | 4 +-
packages/sync/src/repositories/rows.ts | 3 +-
packages/sync/src/services/acquisition.ts | 3 +
packages/sync/src/services/management.ts | 5 +-
packages/sync/test/acquisition-retry.test.ts | 14 +-
packages/sync/test/connector-failure.test.ts | 4 +-
packages/sync/test/package-consumer.ts | 5 +-
packages/sync/test/persistence.test.ts | 2 +-
packages/sync/test/polling-migration.test.ts | 110 +++++++++++++++
packages/sync/test/run-due.test.ts | 132 ++++++++++++++++--
packages/sync/test/shared-backoff.test.ts | 2 +-
packages/sync/test/source-http-error.test.ts | 4 +-
27 files changed, 316 insertions(+), 97 deletions(-)
create mode 100644 packages/sync/test/polling-migration.test.ts
diff --git a/README.md b/README.md
index 31f41fc..cb18955 100644
--- a/README.md
+++ b/README.md
@@ -60,9 +60,9 @@ For a compiled binary, follow the [sample entrypoint](apps/web/backend/src/entry
`splitting: true`. This keeps the temporary cron process from loading the server modules. It asks
the running app to process due syncs, waits for completion, and exits.
-Use `sync.providers` to connect accounts. Create a destination with `sync.api.createDestination()`,
-then link a source to it with `sync.api.createInstallation()`. Supply their `config` values,
-a provider `connection` when needed, and `intervalMs` for the sync schedule. Pass the acting user's
+Use `sync.providers` to connect accounts, then create a sync with `sync.api.createSync()`.
+Supply the source `definition`, its `config`, the `destination`, and a provider `connection`
+when needed. All enabled syncs poll on the shared cron schedule. Pass the acting user's
`actorId` and data owner's `ownerId` from your host's authentication.
See the [sample host](apps/web/backend/src/main.ts),
diff --git a/apps/web/backend/src/services/dashboard/service.ts b/apps/web/backend/src/services/dashboard/service.ts
index 0790683..1800b4f 100644
--- a/apps/web/backend/src/services/dashboard/service.ts
+++ b/apps/web/backend/src/services/dashboard/service.ts
@@ -4,7 +4,6 @@ import type { JsonObject } from '@context-use/open-sync/json';
import { BadRequestError } from '#backend/lib/errors.ts';
-const pollIntervalMs = 900_000;
export type CreateSync = {
source: string;
config?: JsonObject;
@@ -53,7 +52,6 @@ export class DashboardService {
destination: input.destination,
config,
connection,
- intervalMs: pollIntervalMs,
enabled: !definition.provider || !!connection,
});
// Authorization may finish in another tab between the status check and persistence.
diff --git a/apps/web/frontend/src/routes/_workspace/syncs.$id.tsx b/apps/web/frontend/src/routes/_workspace/syncs.$id.tsx
index 5c24168..26289f2 100644
--- a/apps/web/frontend/src/routes/_workspace/syncs.$id.tsx
+++ b/apps/web/frontend/src/routes/_workspace/syncs.$id.tsx
@@ -16,7 +16,6 @@ import {
import { SyncAccount } from './-syncs/sync-account';
import { PendingDeliverables, PollHistory, ReceivedDeliverables } from './-syncs/sync-activity';
-const millisecondsPerMinute = 60_000;
export const Route = createFileRoute('/_workspace/syncs/$id')({
validateSearch: (
search: Record,
@@ -145,10 +144,6 @@ function SyncDetail() {
- Status
- {sync.status.replaceAll('_', ' ')}
- - Sync interval
- - {sync.intervalMs / millisecondsPerMinute} minutes
- - Next scheduled run
- - {nextRun(sync)}
diff --git a/packages/sync/src/db/client.ts b/packages/sync/src/db/client.ts
index 4e53f7b..0d83089 100644
--- a/packages/sync/src/db/client.ts
+++ b/packages/sync/src/db/client.ts
@@ -5,7 +5,18 @@ export function openDatabase(path: string): Database {
const db = new Database(path, { create: true, strict: true });
try {
db.exec('PRAGMA foreign_keys=ON; PRAGMA journal_mode=WAL; PRAGMA busy_timeout=5000;');
- db.transaction(() => db.exec(schema)).immediate();
+ db.transaction(() => {
+ db.exec(schema);
+ const columns = db.query<{ name: string }, []>('PRAGMA table_info(syncs)').all();
+ if (columns.some(({ name }) => name === 'interval_ms')) {
+ // Execute separately so an intermediate failure aborts the entire migration.
+ db.query('ALTER TABLE syncs ADD COLUMN retry_at INTEGER').run();
+ db.query(`UPDATE syncs SET retry_at=next_due_at
+ WHERE status IN ('retrying','interrupted','waiting_for_capacity')`).run();
+ db.query('ALTER TABLE syncs DROP COLUMN interval_ms').run();
+ db.query('ALTER TABLE syncs DROP COLUMN next_due_at').run();
+ }
+ }).immediate();
return db;
} catch (error) {
db.close();
diff --git a/packages/sync/src/db/schema.sql b/packages/sync/src/db/schema.sql
index 7a288f1..469ce53 100644
--- a/packages/sync/src/db/schema.sql
+++ b/packages/sync/src/db/schema.sql
@@ -3,7 +3,7 @@ CREATE TABLE IF NOT EXISTS syncs (
connection TEXT, config TEXT NOT NULL, destination_type TEXT NOT NULL, destination_config TEXT NOT NULL,
enabled INTEGER NOT NULL,
checkpoint TEXT NOT NULL,
- interval_ms INTEGER NOT NULL, next_due_at INTEGER NOT NULL, status TEXT NOT NULL,
+ retry_at INTEGER, status TEXT NOT NULL,
error_code TEXT,
generation INTEGER NOT NULL DEFAULT 0, expires_at INTEGER,
failure_count INTEGER NOT NULL DEFAULT 0, resync INTEGER NOT NULL DEFAULT 0,
diff --git a/packages/sync/src/execution/worker.ts b/packages/sync/src/execution/worker.ts
index b2e7fe1..5df328a 100644
--- a/packages/sync/src/execution/worker.ts
+++ b/packages/sync/src/execution/worker.ts
@@ -27,6 +27,7 @@ export class Worker implements WorkerControl {
#capacityReleased = false;
#closing?: Promise;
#draining?: Promise;
+ #polling = false;
#drainRequested = false;
#resume?: () => void;
constructor(
@@ -55,7 +56,7 @@ export class Worker implements WorkerControl {
if ((!this.#started && !this.#draining) || this.#lifetime.signal.aborted) {
return;
}
- void this.runDue().catch(() => this.input.log({ code: 'runtime_failed' }));
+ void this.requestDrain().catch(() => this.input.log({ code: 'runtime_failed' }));
}
/** A deterministic dispatch round for hosts/tests that do not start the background worker. */
tick(): Promise {
@@ -65,9 +66,16 @@ export class Worker implements WorkerControl {
this.cleanup(),
);
}
- /** Drain runnable pages and deliveries, sharing one run across concurrent host triggers. */
+ /** Poll enabled syncs and drain runnable work, sharing a round across concurrent host triggers. */
runDue(): Promise {
this.ensureOpen();
+ if (!this.#polling) {
+ this.input.acquisition.poll();
+ this.#polling = true;
+ }
+ return this.requestDrain();
+ }
+ private requestDrain(): Promise {
this.#drainRequested = true;
this.#resume?.();
this.#draining ??= Promise.resolve().then(() => this.drain());
@@ -98,6 +106,7 @@ export class Worker implements WorkerControl {
} finally {
this.#resume = undefined;
this.#draining = undefined;
+ this.#polling = false;
}
}
private dispatch(): void {
diff --git a/packages/sync/src/http/controller.ts b/packages/sync/src/http/controller.ts
index 6177455..16e0bc3 100644
--- a/packages/sync/src/http/controller.ts
+++ b/packages/sync/src/http/controller.ts
@@ -15,7 +15,6 @@ const syncBody = t.Object({
config,
destination: t.Object({ type: identifier, input: config }),
connection: t.Optional(t.Object({ id: identifier, service: identifier })),
- intervalMs: t.Optional(t.Integer({ minimum: 1 })),
enabled: t.Optional(t.Boolean()),
});
const resourceParams = { params: t.Object({ id: identifier }) };
diff --git a/packages/sync/src/models/definition.ts b/packages/sync/src/models/definition.ts
index fab0ddc..4893b5f 100644
--- a/packages/sync/src/models/definition.ts
+++ b/packages/sync/src/models/definition.ts
@@ -54,7 +54,7 @@ export interface SyncStep {
* Include any cross-iteration update boundary as well as the current page position.
*/
checkpoint: JsonValue;
- /** False continues this iteration immediately; true schedules the next poll.
+ /** False continues this iteration immediately; true waits for the shared cron or a manual run.
* Completion retains the returned checkpoint; the source clears any exhausted page cursor.
*/
complete: boolean;
diff --git a/packages/sync/src/models/sync.ts b/packages/sync/src/models/sync.ts
index b25e5e8..018b6da 100644
--- a/packages/sync/src/models/sync.ts
+++ b/packages/sync/src/models/sync.ts
@@ -11,8 +11,7 @@ export interface Sync {
destination: { type: string; config: JsonObject };
enabled: boolean;
checkpoint: JsonValue;
- intervalMs: number;
- nextDueAt: number;
+ retryAt: number | null;
status: SyncStatus;
errorCode: string | null;
}
@@ -21,7 +20,6 @@ export interface CreateSync extends Scope {
connection?: ConnectionRef;
config: JsonObject;
destination: { type: string; input: JsonObject };
- intervalMs?: number;
enabled?: boolean;
}
@@ -51,8 +49,7 @@ export interface SyncSummary {
destinationType: string;
connection?: ConnectionRef;
enabled: boolean;
- intervalMs: number;
- nextDueAt: number;
+ retryAt: number | null;
status: SyncStatus;
errorCode: string | null;
}
@@ -63,8 +60,7 @@ export function summarizeSync(sync: Sync): SyncSummary {
destinationType: sync.destination.type,
connection: sync.connection,
enabled: sync.enabled,
- intervalMs: sync.intervalMs,
- nextDueAt: sync.nextDueAt,
+ retryAt: sync.retryAt,
status: sync.status,
errorCode: sync.errorCode,
};
diff --git a/packages/sync/src/repositories/acquisition/contract.ts b/packages/sync/src/repositories/acquisition/contract.ts
index 29c3450..e65952e 100644
--- a/packages/sync/src/repositories/acquisition/contract.ts
+++ b/packages/sync/src/repositories/acquisition/contract.ts
@@ -9,6 +9,7 @@ export interface AcquisitionLease extends Scope {
failureCount: number;
}
export interface AcquisitionRepository {
+ poll(): void;
capacityReleased(): void;
claim(leaseMs: number): AcquisitionLease | undefined;
hasCapacity(): boolean;
diff --git a/packages/sync/src/repositories/acquisition/lease.ts b/packages/sync/src/repositories/acquisition/lease.ts
index a561daf..a7c230b 100644
--- a/packages/sync/src/repositories/acquisition/lease.ts
+++ b/packages/sync/src/repositories/acquisition/lease.ts
@@ -30,8 +30,8 @@ export function claimAcquisition({
WHERE completed_at IS NULL AND EXISTS (SELECT 1 FROM syncs
WHERE syncs.owner_id=sync_polls.owner_id AND syncs.id=sync_polls.sync_id
AND enabled=1 AND expires_at<=?)`).run(now);
- db.query(`UPDATE syncs SET next_due_at=?,status='interrupted',error_code='lease_expired',expires_at=NULL
- WHERE enabled=1 AND expires_at<=?`).run(now, now);
+ db.query(`UPDATE syncs SET retry_at=NULL,status='interrupted',error_code='lease_expired',expires_at=NULL
+ WHERE enabled=1 AND expires_at<=?`).run(now);
const row = db
.query<
{
@@ -43,8 +43,8 @@ export function claimAcquisition({
},
[number]
>(`SELECT id,owner_id,generation,resync,failure_count FROM syncs
- WHERE enabled=1 AND expires_at IS NULL AND next_due_at<=?
- ORDER BY next_due_at,generation,id LIMIT 1`)
+ WHERE enabled=1 AND status!='succeeded' AND expires_at IS NULL AND COALESCE(retry_at,0)<=?
+ ORDER BY COALESCE(retry_at,0),generation,id LIMIT 1`)
.get(now);
if (!row) {
return;
@@ -88,13 +88,13 @@ export function finishAcquisition(
lease.ownerId,
lease.sync.id,
);
- db.query(`UPDATE syncs SET status=?,error_code=?,next_due_at=?,expires_at=NULL,
+ db.query(`UPDATE syncs SET status=?,error_code=?,retry_at=?,expires_at=NULL,
enabled=CASE WHEN ? THEN 0 ELSE enabled END,
failure_count=COALESCE(?,failure_count),resync=CASE WHEN ? THEN 0 ELSE resync END
WHERE owner_id=? AND id=?`).run(
input.pause ? 'disabled' : input.state,
input.errorCode ?? null,
- now + input.delay,
+ input.delay > 0 ? now + input.delay : null,
Number(input.pause ?? false),
input.failureCount ?? null,
Number(input.state === 'succeeded'),
diff --git a/packages/sync/src/repositories/acquisition/sqlite.ts b/packages/sync/src/repositories/acquisition/sqlite.ts
index 7412814..7a72447 100644
--- a/packages/sync/src/repositories/acquisition/sqlite.ts
+++ b/packages/sync/src/repositories/acquisition/sqlite.ts
@@ -15,10 +15,15 @@ import { writeRecord } from './records';
export class SqliteAcquisition implements AcquisitionRepository {
constructor(private readonly input: { db: Database; limits: QueueLimits }) {}
+ poll(): void {
+ this.input.db
+ .query("UPDATE syncs SET status='ready' WHERE enabled=1 AND status='succeeded'")
+ .run();
+ }
capacityReleased(): void {
this.input.db
- .query("UPDATE syncs SET next_due_at=? WHERE enabled=1 AND status='waiting_for_capacity'")
- .run(Date.now());
+ .query("UPDATE syncs SET retry_at=NULL WHERE enabled=1 AND status='waiting_for_capacity'")
+ .run();
}
claim(leaseMs: number) {
return claimAcquisition({ ...this.input, leaseMs });
@@ -52,7 +57,7 @@ export class SqliteAcquisition implements AcquisitionRepository {
db,
lease: input.lease,
state: page.complete ? 'succeeded' : 'ready',
- delay: page.complete ? sync.intervalMs : 0,
+ delay: 0,
failureCount: page.complete ? 0 : undefined,
});
}).immediate();
diff --git a/packages/sync/src/repositories/assets/sqlite.ts b/packages/sync/src/repositories/assets/sqlite.ts
index 8c3b245..f4d15c9 100644
--- a/packages/sync/src/repositories/assets/sqlite.ts
+++ b/packages/sync/src/repositories/assets/sqlite.ts
@@ -103,9 +103,9 @@ export class SqliteAssets implements AssetRepository {
if (deleted) {
this.input.db
.query(
- "UPDATE syncs SET next_due_at=? WHERE enabled=1 AND status='waiting_for_capacity' AND NOT (owner_id=? AND id=?)",
+ "UPDATE syncs SET retry_at=NULL WHERE enabled=1 AND status='waiting_for_capacity' AND NOT (owner_id=? AND id=?)",
)
- .run(Date.now(), deleted.owner_id, deleted.sync_id);
+ .run(deleted.owner_id, deleted.sync_id);
}
}
garbage(): string[] {
diff --git a/packages/sync/src/repositories/catalog/sqlite.ts b/packages/sync/src/repositories/catalog/sqlite.ts
index b34a70a..8ab2280 100644
--- a/packages/sync/src/repositories/catalog/sqlite.ts
+++ b/packages/sync/src/repositories/catalog/sqlite.ts
@@ -7,8 +7,6 @@ import type { CreateSync, SyncPoll } from '../../models/sync';
import { readSync } from '../rows';
import type { CatalogRepository } from './contract';
-const defaultIntervalMs = 60_000;
-
export class SqliteCatalog implements CatalogRepository {
constructor(private readonly db: Database) {}
createSync(
@@ -19,8 +17,8 @@ export class SqliteCatalog implements CatalogRepository {
) {
const id = `sync_${crypto.randomUUID()}`;
this.db
- .query(`INSERT INTO syncs (owner_id,id,definition_id,connection,config,destination_type,destination_config,enabled,checkpoint,interval_ms,next_due_at,status)
- VALUES (?,?,?,?,?,?,?,?,?,?,?,?)`)
+ .query(`INSERT INTO syncs (owner_id,id,definition_id,connection,config,destination_type,destination_config,enabled,checkpoint,status)
+ VALUES (?,?,?,?,?,?,?,?,?,?)`)
.run(
input.ownerId,
id,
@@ -31,8 +29,6 @@ export class SqliteCatalog implements CatalogRepository {
canonicalJson(input.destination.config).json,
Number(input.enabled ?? true),
canonicalJson(input.initialCheckpoint).json,
- input.intervalMs ?? defaultIntervalMs,
- Date.now(),
input.enabled === false ? 'disabled' : 'ready',
);
return this.sync({ ...input, id });
@@ -69,8 +65,8 @@ export class SqliteCatalog implements CatalogRepository {
}
this.db
.query(`UPDATE syncs SET connection=?,enabled=1,
- status='ready',error_code=NULL,next_due_at=? WHERE owner_id=? AND id=?`)
- .run(canonicalJson(input.connection).json, Date.now(), input.ownerId, input.id);
+ status='ready',error_code=NULL,retry_at=NULL WHERE owner_id=? AND id=?`)
+ .run(canonicalJson(input.connection).json, input.ownerId, input.id);
return this.sync(input);
})
.immediate();
@@ -88,12 +84,11 @@ export class SqliteCatalog implements CatalogRepository {
.run(input.enabled ? 'ready' : 'disabled', input.ownerId, input.id);
this.db
.query(
- `UPDATE syncs SET enabled=?,status=?,error_code=NULL,next_due_at=?,generation=generation+1,expires_at=NULL WHERE owner_id=? AND id=?`,
+ `UPDATE syncs SET enabled=?,status=?,error_code=NULL,retry_at=NULL,generation=generation+1,expires_at=NULL WHERE owner_id=? AND id=?`,
)
.run(
Number(input.enabled),
input.enabled ? 'ready' : 'disabled',
- Date.now(),
input.ownerId,
input.id,
);
@@ -115,8 +110,8 @@ export class SqliteCatalog implements CatalogRepository {
fail('busy');
}
this.db
- .query(`UPDATE syncs SET next_due_at=? WHERE owner_id=? AND id=?`)
- .run(Date.now(), input.ownerId, input.id);
+ .query(`UPDATE syncs SET status='ready',retry_at=NULL WHERE owner_id=? AND id=?`)
+ .run(input.ownerId, input.id);
})
.immediate();
}
@@ -132,9 +127,9 @@ export class SqliteCatalog implements CatalogRepository {
.run(Date.now(), input.ownerId, input.id);
this.db
.query(
- `UPDATE syncs SET checkpoint=?,status='ready',error_code=NULL,next_due_at=?,resync=1,failure_count=0,generation=generation+1,expires_at=NULL WHERE owner_id=? AND id=?`,
+ `UPDATE syncs SET checkpoint=?,status='ready',error_code=NULL,retry_at=NULL,resync=1,failure_count=0,generation=generation+1,expires_at=NULL WHERE owner_id=? AND id=?`,
)
- .run(canonicalJson(input.checkpoint).json, Date.now(), input.ownerId, input.id);
+ .run(canonicalJson(input.checkpoint).json, input.ownerId, input.id);
})
.immediate();
}
@@ -150,8 +145,8 @@ export class SqliteCatalog implements CatalogRepository {
.run(input.ownerId, input.id);
this.db.query('DELETE FROM syncs WHERE owner_id=? AND id=?').run(input.ownerId, input.id);
this.db
- .query(`UPDATE syncs SET next_due_at=? WHERE enabled=1 AND status='waiting_for_capacity'`)
- .run(Date.now());
+ .query(`UPDATE syncs SET retry_at=NULL WHERE enabled=1 AND status='waiting_for_capacity'`)
+ .run();
})
.immediate();
}
diff --git a/packages/sync/src/repositories/delivery/sqlite.ts b/packages/sync/src/repositories/delivery/sqlite.ts
index 926a4d7..798f90b 100644
--- a/packages/sync/src/repositories/delivery/sqlite.ts
+++ b/packages/sync/src/repositories/delivery/sqlite.ts
@@ -59,9 +59,9 @@ export class SqliteDeliveries implements DeliveryRepository {
// Capacity is global; all paused acquisitions can compete for the released budget.
this.db
.query(
- "UPDATE syncs SET next_due_at=? WHERE enabled=1 AND status='waiting_for_capacity'",
+ "UPDATE syncs SET retry_at=NULL WHERE enabled=1 AND status='waiting_for_capacity'",
)
- .run(Date.now());
+ .run();
} else {
this.db
.query(
diff --git a/packages/sync/src/repositories/rows.ts b/packages/sync/src/repositories/rows.ts
index 089bdc7..0eb89c1 100644
--- a/packages/sync/src/repositories/rows.ts
+++ b/packages/sync/src/repositories/rows.ts
@@ -23,8 +23,7 @@ export function readSync(input: { db: Database; scope: Resource }): Sync {
},
enabled: row.enabled === 1,
checkpoint: JSON.parse(String(row.checkpoint)),
- intervalMs: Number(row.interval_ms),
- nextDueAt: Number(row.next_due_at),
+ retryAt: row.retry_at === null ? null : Number(row.retry_at),
status: row.status as SyncStatus,
errorCode: row.error_code === null ? null : String(row.error_code),
};
diff --git a/packages/sync/src/services/acquisition.ts b/packages/sync/src/services/acquisition.ts
index be4ff1f..f97c78c 100644
--- a/packages/sync/src/services/acquisition.ts
+++ b/packages/sync/src/services/acquisition.ts
@@ -24,6 +24,9 @@ export class AcquisitionService {
log: Logger;
},
) {}
+ poll() {
+ this.input.repository.poll();
+ }
capacityReleased() {
this.input.repository.capacityReleased();
}
diff --git a/packages/sync/src/services/management.ts b/packages/sync/src/services/management.ts
index 9cb287c..d47c018 100644
--- a/packages/sync/src/services/management.ts
+++ b/packages/sync/src/services/management.ts
@@ -4,7 +4,7 @@ import type { ConnectionRef } from '../models/definition';
import { fail } from '../models/error';
import type { Resource, Scope } from '../models/identity';
import type { JsonObject } from '../models/json';
-import { positive, type QueueLimits } from '../models/limits';
+import type { QueueLimits } from '../models/limits';
import type { Registry } from '../models/registry';
import { type CreateSync, summarizeSync } from '../models/sync';
import { identifier, validate } from '../models/validation';
@@ -67,9 +67,6 @@ export class SyncManagement {
this.guard(input);
const { definition } = this.input.registry.definition(input.definition);
const config = validate({ value: input.config, schema: definition.configSchema }) as JsonObject;
- if (input.intervalMs !== undefined) {
- positive(input.intervalMs);
- }
if (definition.provider && (input.connection || input.enabled !== false)) {
await bindProvider({
actorId: input.actorId,
diff --git a/packages/sync/test/acquisition-retry.test.ts b/packages/sync/test/acquisition-retry.test.ts
index e8b0201..5646faa 100644
--- a/packages/sync/test/acquisition-retry.test.ts
+++ b/packages/sync/test/acquisition-retry.test.ts
@@ -132,7 +132,7 @@ test('source failures back off durably despite partial progress, then reset on s
checkpoint: 1,
status: 'retrying',
errorCode: 'execution_failed',
- nextDueAt: now + delay,
+ retryAt: now + delay,
});
expect(events.at(-1)?.fields?.retryAfterMs).toBe(delay);
const attempts = events.length;
@@ -160,7 +160,7 @@ test('source failures back off durably despite partial progress, then reset on s
await engine.tick();
await engine.tick();
await engine.tick();
- expect(engine.api.sync({ ...beta, id: other.id }).nextDueAt).toBe(now + retryMs);
+ expect(engine.api.sync({ ...beta, id: other.id }).retryAt).toBe(now + retryMs);
await engine.api.setEnabled({ ...beta, id: other.id, enabled: false });
failing = false;
@@ -170,7 +170,7 @@ test('source failures back off durably despite partial progress, then reset on s
failing = true;
engine.api.runNow(scope);
await engine.tick();
- expect(engine.api.sync(scope).nextDueAt).toBe(now + retryMs);
+ expect(engine.api.sync(scope).retryAt).toBe(now + retryMs);
expect(events.at(-1)?.fields?.failureCount).toBe(1);
} finally {
await engine.close();
@@ -226,26 +226,26 @@ test.each(['records', 'assets'])(
const sync = await configure(engine);
const scope = { ...alpha, id: sync.id };
await engine.tick();
- expect(engine.api.sync(scope).nextDueAt).toBe(now + retryMs);
+ expect(engine.api.sync(scope).retryAt).toBe(now + retryMs);
expect(events.at(-1)?.fields).toMatchObject({ httpStatus: 429, retryAfterMs: retryMs });
expect(JSON.stringify(events)).not.toContain('private upstream payload');
now += retryMs;
await engine.tick();
const secondDelay = 60_000;
- expect(engine.api.sync(scope).nextDueAt).toBe(now + secondDelay);
+ expect(engine.api.sync(scope).retryAt).toBe(now + secondDelay);
now += secondDelay;
failing = false;
await engine.tick();
expect(savedSync({ path: files.path, scope: scope })).toMatchObject({
status: 'ready',
- nextDueAt: now,
+ retryAt: null,
});
await engine.close();
engine = createSyncRuntime(options);
failing = true;
await engine.tick();
const thirdDelay = 120_000;
- expect(engine.api.sync(scope).nextDueAt).toBe(now + thirdDelay);
+ expect(engine.api.sync(scope).retryAt).toBe(now + thirdDelay);
} finally {
await engine.close();
clock.mockRestore();
diff --git a/packages/sync/test/connector-failure.test.ts b/packages/sync/test/connector-failure.test.ts
index 0ed31cc..fefe882 100644
--- a/packages/sync/test/connector-failure.test.ts
+++ b/packages/sync/test/connector-failure.test.ts
@@ -148,7 +148,7 @@ test('connector failures use engine backoff despite timing headers and retain on
const secondDelay = 60_000;
expect(savedSync({ path: files.path, scope: scope })).toMatchObject({
checkpoint: 0,
- nextDueAt: now + firstDelay,
+ retryAt: now + firstDelay,
});
now += firstDelay;
const providerDelay = 120_000;
@@ -156,7 +156,7 @@ test('connector failures use engine backoff despite timing headers and retain on
await engine.tick();
expect(savedSync({ path: files.path, scope: scope })).toMatchObject({
checkpoint: 0,
- nextDueAt: now + secondDelay,
+ retryAt: now + secondDelay,
});
expect(events.at(-1)?.fields).toMatchObject({
httpStatus: rateLimited,
diff --git a/packages/sync/test/package-consumer.ts b/packages/sync/test/package-consumer.ts
index 88a5976..1d3752f 100644
--- a/packages/sync/test/package-consumer.ts
+++ b/packages/sync/test/package-consumer.ts
@@ -14,7 +14,7 @@ if (await runOpenSyncCron()) {
// Keep the installed provider runtime and credential persistence real; simulate only GitHub.
const providerFetch = globalThis.fetch;
const okStatus = 200;
-const AFTER_INTERVAL_MS = 150;
+const IDLE_OBSERVATION_MS = 150;
globalThis.fetch = Object.assign((input: RequestInfo | URL) => {
const url = input instanceof Request ? input.url : String(input);
if (url !== 'https://api.github.com/user') {
@@ -177,7 +177,6 @@ try {
config: {},
definition: definition.definition.id,
connection: { id: connection.id, service: connection.service },
- intervalMs: 100,
});
// Exercise the installed API with a private crontab subprocess, including a compiled host
// launched from outside its source tree. No system crontab or public server is involved.
@@ -187,7 +186,7 @@ try {
throw new Error('Independent consumer did not receive its record');
}
const pollCount = sync.api.polls({ ...scope, id: created.id }).polls.length;
- await Bun.sleep(AFTER_INTERVAL_MS);
+ await Bun.sleep(IDLE_OBSERVATION_MS);
if (sync.api.polls({ ...scope, id: created.id }).polls.length !== pollCount) {
throw new Error('An in-process timer bypassed the cron schedule');
}
diff --git a/packages/sync/test/persistence.test.ts b/packages/sync/test/persistence.test.ts
index 1ae50cb..efb9ef8 100644
--- a/packages/sync/test/persistence.test.ts
+++ b/packages/sync/test/persistence.test.ts
@@ -13,7 +13,7 @@ test.each([false, true])(
const before = f.catalog.sync({ ...alpha, id: f.sync.id });
const pollsBefore = f.catalog.polls({ ...alpha, id: f.sync.id });
f.db.exec(
- "CREATE TRIGGER fail_checkpoint BEFORE UPDATE OF next_due_at ON syncs BEGIN SELECT RAISE(ABORT,'injected'); END;",
+ "CREATE TRIGGER fail_checkpoint BEFORE UPDATE OF retry_at ON syncs BEGIN SELECT RAISE(ABORT,'injected'); END;",
);
expect(() =>
f.acquisition.commit({
diff --git a/packages/sync/test/polling-migration.test.ts b/packages/sync/test/polling-migration.test.ts
new file mode 100644
index 0000000..3838842
--- /dev/null
+++ b/packages/sync/test/polling-migration.test.ts
@@ -0,0 +1,110 @@
+import { Database } from 'bun:sqlite';
+import { expect, test } from 'bun:test';
+import { openDatabase } from '../src/db/client';
+import { defaultLimits } from '../src/models/limits';
+import { SqliteAcquisition } from '../src/repositories/acquisition/sqlite';
+import { SqliteCatalog } from '../src/repositories/catalog/sqlite';
+import { alpha, fixture, storage } from './support';
+
+const legacySchema = `CREATE TABLE syncs (
+ owner_id TEXT NOT NULL, id TEXT NOT NULL, definition_id TEXT NOT NULL,
+ connection TEXT, config TEXT NOT NULL, destination_type TEXT NOT NULL, destination_config TEXT NOT NULL,
+ enabled INTEGER NOT NULL, checkpoint TEXT NOT NULL,
+ interval_ms INTEGER NOT NULL, next_due_at INTEGER NOT NULL, status TEXT NOT NULL,
+ error_code TEXT, generation INTEGER NOT NULL DEFAULT 0, expires_at INTEGER,
+ failure_count INTEGER NOT NULL DEFAULT 0, resync INTEGER NOT NULL DEFAULT 0,
+ PRIMARY KEY(owner_id,id)
+);`;
+const leaseMs = 60_000;
+const futureRetry = 9_000_000_000_000;
+function legacyDatabase(path: string) {
+ const db = new Database(path);
+ db.exec(legacySchema);
+ for (const state of ['succeeded', 'retrying', 'disabled', 'ready', 'running']) {
+ db.query(`INSERT INTO syncs (owner_id,id,definition_id,config,destination_type,
+ destination_config,enabled,checkpoint,interval_ms,next_due_at,status,generation,expires_at,
+ failure_count,resync) VALUES (?,?,?,'{"count":3}','local','{}',?, '2',60000,?,?,7,?,2,1)`).run(
+ alpha.ownerId,
+ state,
+ fixture.definition.id,
+ Number(state !== 'disabled'),
+ futureRetry,
+ state,
+ state === 'running' ? futureRetry : null,
+ );
+ }
+ return db;
+}
+
+test('legacy polling schedules migrate without losing progress, leases, backoff or queued data', () => {
+ const files = storage();
+ let db = legacyDatabase(files.path);
+ // Real foreign keys and populated dependent tables must survive the column removal.
+ db.exec(`CREATE TABLE retained (
+ sync_id TEXT, owner_id TEXT, payload TEXT,
+ FOREIGN KEY(owner_id,sync_id) REFERENCES syncs(owner_id,id)
+ ); INSERT INTO retained VALUES ('succeeded','alpha','queued records and assets');`);
+ const before = db.query('SELECT * FROM syncs ORDER BY id').all() as Record[];
+ db.close();
+ try {
+ db = openDatabase(files.path);
+ const after = db.query('SELECT * FROM syncs ORDER BY id').all();
+ expect(after).toEqual(
+ before.map(({ interval_ms: _interval, next_due_at, ...row }) => ({
+ ...row,
+ retry_at: row.status === 'retrying' ? next_due_at : null,
+ })),
+ );
+ expect(db.query('PRAGMA foreign_key_check').all()).toEqual([]);
+ expect(db.query('SELECT * FROM retained').all()).toEqual([
+ { sync_id: 'succeeded', owner_id: 'alpha', payload: 'queued records and assets' },
+ ]);
+ db.close();
+ db = openDatabase(files.path);
+ expect(db.query('SELECT * FROM syncs ORDER BY id').all()).toEqual(after);
+ const catalog = new SqliteCatalog(db);
+ const acquisition = new SqliteAcquisition({ db, limits: defaultLimits });
+ expect(acquisition.claim(leaseMs)?.sync.id).toBe('ready');
+ expect(acquisition.claim(leaseMs)).toBeUndefined();
+ acquisition.poll();
+ expect(acquisition.claim(leaseMs)?.sync.id).toBe('succeeded');
+ expect(acquisition.claim(leaseMs)).toBeUndefined();
+ expect(
+ catalog.createSync({
+ ...alpha,
+ definition: 'test',
+ config: { count: 1 },
+ destination: { type: 'local', config: {} },
+ initialCheckpoint: 0,
+ }).retryAt,
+ ).toBeNull();
+ } finally {
+ db.close();
+ files.close();
+ }
+});
+
+test('failed migration rolls back all schema and data changes and can be retried', () => {
+ const files = storage();
+ let db = legacyDatabase(files.path);
+ db.exec(
+ "CREATE TRIGGER fail_migration BEFORE UPDATE ON syncs BEGIN SELECT RAISE(ABORT,'injected'); END",
+ );
+ const before = db.query('SELECT * FROM syncs ORDER BY id').all();
+ db.close();
+ try {
+ expect(() => openDatabase(files.path)).toThrow();
+ db = new Database(files.path);
+ expect(db.query('SELECT * FROM syncs ORDER BY id').all()).toEqual(before);
+ expect(db.query("SELECT name FROM pragma_table_info('syncs')").all()).not.toContainEqual({
+ name: 'retry_at',
+ });
+ db.exec('DROP TRIGGER fail_migration');
+ db.close();
+ db = openDatabase(files.path);
+ expect(db.query('SELECT id FROM syncs').all()).toHaveLength(before.length);
+ } finally {
+ db.close();
+ files.close();
+ }
+});
diff --git a/packages/sync/test/run-due.test.ts b/packages/sync/test/run-due.test.ts
index 0f3f792..66091b0 100644
--- a/packages/sync/test/run-due.test.ts
+++ b/packages/sync/test/run-due.test.ts
@@ -1,11 +1,10 @@
-import { expect, test } from 'bun:test';
+import { expect, spyOn, test } from 'bun:test';
import { createSyncRuntime } from '../src/runtime';
import { isolatedScheduler, runRegisteredCron } from './cron-support';
import { accepted, alpha, beta, configure, fixture, runtime, storage } from './support';
const RECORD_COUNT = 3;
-const POLL_INTERVAL_MS = 100;
-const AFTER_INTERVAL_MS = 150;
+const IDLE_OBSERVATION_MS = 150;
test('one host trigger drains all due syncs and deliveries, coalesces overlap, and respects pauses', async () => {
const files = storage();
@@ -33,7 +32,6 @@ test('one host trigger drains all due syncs and deliveries, coalesces overlap, a
definition: 'test',
config: { count: RECORD_COUNT },
destination: { type: 'local', input: {} },
- intervalMs: POLL_INTERVAL_MS,
}),
})),
);
@@ -45,8 +43,7 @@ test('one host trigger drains all due syncs and deliveries, coalesces overlap, a
expect(engine.api.status(scope).queue.pendingRecords).toBe(0);
expect(engine.api.polls({ ...scope, id: sync.id }).polls).toHaveLength(1);
}
- // A host may choose its own timer or external scheduler without starting background work.
- await Bun.sleep(AFTER_INTERVAL_MS);
+ // Every host trigger polls again, without waiting for a per-sync deadline.
const [paused, active] = syncs;
await engine.api.setEnabled({ ...paused!.scope, id: paused!.sync.id, enabled: false });
expect(engine.api.polls({ ...active!.scope, id: active!.sync.id }).polls).toHaveLength(1);
@@ -83,7 +80,6 @@ test('cron joins startup work and remains the only automatic polling trigger', a
definition: 'test',
config: { count: 1 },
destination: { type: 'local', input: {} },
- intervalMs: POLL_INTERVAL_MS,
});
await f.engine.start();
await entered.promise;
@@ -92,7 +88,7 @@ test('cron joins startup work and remains the only automatic polling trigger', a
release.resolve();
await running;
expect(steps).toBe(1);
- await Bun.sleep(AFTER_INTERVAL_MS);
+ await Bun.sleep(IDLE_OBSERVATION_MS);
expect(steps).toBe(1);
await runRegisteredCron(table.table);
expect(steps).toBe(2);
@@ -129,3 +125,123 @@ test('shutdown interrupts a host trigger waiting for due work', async () => {
await f.close();
}
});
+
+test('manual wake-ups leave completed syncs idle and overlapping cron triggers share one poll', async () => {
+ await using table = await isolatedScheduler();
+ const blocked = Promise.withResolvers();
+ const release = Promise.withResolvers();
+ const calls: number[] = [];
+ const f = runtime({
+ registration: {
+ ...fixture,
+ load: () => ({
+ async step({ config }) {
+ calls.push(Number(config.count));
+ if (config.count === 2) {
+ blocked.resolve();
+ await release.promise;
+ }
+ return { records: [], checkpoint: 0, complete: true };
+ },
+ }),
+ },
+ });
+ const create = (count: number) =>
+ f.engine.api.createSync({
+ ...alpha,
+ definition: 'test',
+ config: { count },
+ destination: { type: 'local', input: {} },
+ });
+ try {
+ const first = await create(1);
+ await f.engine.runDue();
+ await f.engine.start();
+ await create(2);
+ await blocked.promise;
+ expect(calls).toEqual([1, 2]);
+ const cron = f.engine.runDue();
+ // Starting a cron round during a manual drain still includes the previously completed sync.
+ // Wait for its committed poll while the other sync remains in flight.
+ await until(
+ () =>
+ f.engine.api.polls({ ...alpha, id: first.id }).polls.length === 2 &&
+ f.engine.api.sync({ ...alpha, id: first.id }).status === 'succeeded',
+ );
+ expect(f.engine.runDue()).toBe(cron);
+ release.resolve();
+ await cron;
+ expect(calls).toEqual([1, 2, 1]);
+ const afterManualPolls = 3;
+ f.engine.api.runNow({ ...alpha, id: first.id });
+ await until(
+ () =>
+ f.engine.api.polls({ ...alpha, id: first.id }).polls.length === afterManualPolls &&
+ f.engine.api.sync({ ...alpha, id: first.id }).status === 'succeeded',
+ );
+ expect(calls).toEqual([1, 2, 1, 1]);
+ await runRegisteredCron(table.table);
+ expect(calls.filter((count) => count === 2)).toHaveLength(2);
+ } finally {
+ release.resolve();
+ await f.close();
+ }
+});
+
+async function until(check: () => boolean) {
+ const timeoutMs = 5_000;
+ const deadline = Date.now() + timeoutMs;
+ while (!check()) {
+ if (Date.now() >= deadline) {
+ throw new Error('Sync did not reach the expected state');
+ }
+ await Bun.sleep(1);
+ }
+}
+
+test('each cron round polls healthy syncs without bypassing source retry backoff', async () => {
+ let now = Date.now();
+ const clock = spyOn(Date, 'now').mockImplementation(() => now);
+ const calls: number[] = [];
+ let failing = true;
+ const f = runtime({
+ registration: {
+ ...fixture,
+ load: () => ({
+ step({ config }) {
+ calls.push(Number(config.count));
+ if (config.count === 2 && failing) {
+ return Promise.reject(new Error('temporary failure'));
+ }
+ return Promise.resolve({ records: [], checkpoint: 0, complete: true });
+ },
+ }),
+ },
+ });
+ try {
+ const syncs = await Promise.all(
+ [1, 2].map((count) =>
+ f.engine.api.createSync({
+ ...alpha,
+ definition: 'test',
+ config: { count },
+ destination: { type: 'local', input: {} },
+ }),
+ ),
+ );
+ await f.engine.runDue();
+ await f.engine.runDue();
+ expect(calls.filter((count) => count === 1)).toHaveLength(2);
+ expect(calls.filter((count) => count === 2)).toHaveLength(1);
+ const scope = { ...alpha, id: syncs[1]!.id };
+ expect(f.engine.api.sync(scope).retryAt).toBe(now + f.options.timing.retryMs);
+ now += f.options.timing.retryMs;
+ failing = false;
+ await f.engine.runDue();
+ expect(calls.filter((count) => count === 2)).toHaveLength(2);
+ expect(f.engine.api.sync(scope)).toMatchObject({ status: 'succeeded', retryAt: null });
+ } finally {
+ await f.close();
+ clock.mockRestore();
+ }
+});
diff --git a/packages/sync/test/shared-backoff.test.ts b/packages/sync/test/shared-backoff.test.ts
index 4a57a7e..9d01dee 100644
--- a/packages/sync/test/shared-backoff.test.ts
+++ b/packages/sync/test/shared-backoff.test.ts
@@ -45,7 +45,7 @@ test('source and delivery failures use the same durable exponential backoff', as
});
for (const delay of delays) {
const due = now + delay;
- expect(engine.api.sync({ ...alpha, id: source.id }).nextDueAt).toBe(due);
+ expect(engine.api.sync({ ...alpha, id: source.id }).retryAt).toBe(due);
expect(
engine.api.deliveries({ ...alpha, syncId: destinationSync.id }).deliveries[0]
?.nextAttemptAt,
diff --git a/packages/sync/test/source-http-error.test.ts b/packages/sync/test/source-http-error.test.ts
index 11be518..592a591 100644
--- a/packages/sync/test/source-http-error.test.ts
+++ b/packages/sync/test/source-http-error.test.ts
@@ -46,7 +46,7 @@ test.each(
checkpoint: 0,
status: 'retrying',
errorCode: `source_http_${status}`,
- nextDueAt: now + retryMs,
+ retryAt: now + retryMs,
});
expect(events[0]?.fields).toEqual({
httpStatus: status,
@@ -174,7 +174,7 @@ test.each([rateLimited, forbidden, unavailable])(
checkpoint: 1,
status: status === forbidden ? 'disabled' : 'retrying',
errorCode: `source_http_${status}`,
- nextDueAt: now + delay,
+ retryAt: now + delay,
});
expect(events.at(-1)?.fields).toEqual({
httpStatus: status,
From a349b093aefda4c07fa64ea7c0a34b5d4953e676 Mon Sep 17 00:00:00 2001
From: massimoalbarello
Date: Fri, 2 Oct 2026 19:31:34 +0100
Subject: [PATCH 2/5] fix: migrate sync databases through versioned history
---
packages/sync/src/db/client.ts | 15 +-
packages/sync/src/db/migrate.ts | 32 +++
.../src/db/migrations/0000-initial-schema.ts | 10 +
.../db/migrations/0001-centralize-polling.ts | 9 +
packages/sync/src/db/schema.sql | 2 +-
packages/sync/test/polling-migration.test.ts | 201 +++++++++++++++---
6 files changed, 220 insertions(+), 49 deletions(-)
create mode 100644 packages/sync/src/db/migrate.ts
create mode 100644 packages/sync/src/db/migrations/0000-initial-schema.ts
create mode 100644 packages/sync/src/db/migrations/0001-centralize-polling.ts
diff --git a/packages/sync/src/db/client.ts b/packages/sync/src/db/client.ts
index 0d83089..9813eb6 100644
--- a/packages/sync/src/db/client.ts
+++ b/packages/sync/src/db/client.ts
@@ -1,22 +1,11 @@
import { Database } from 'bun:sqlite';
-import schema from './schema.sql' with { type: 'text' };
+import { runMigrations } from './migrate';
export function openDatabase(path: string): Database {
const db = new Database(path, { create: true, strict: true });
try {
db.exec('PRAGMA foreign_keys=ON; PRAGMA journal_mode=WAL; PRAGMA busy_timeout=5000;');
- db.transaction(() => {
- db.exec(schema);
- const columns = db.query<{ name: string }, []>('PRAGMA table_info(syncs)').all();
- if (columns.some(({ name }) => name === 'interval_ms')) {
- // Execute separately so an intermediate failure aborts the entire migration.
- db.query('ALTER TABLE syncs ADD COLUMN retry_at INTEGER').run();
- db.query(`UPDATE syncs SET retry_at=next_due_at
- WHERE status IN ('retrying','interrupted','waiting_for_capacity')`).run();
- db.query('ALTER TABLE syncs DROP COLUMN interval_ms').run();
- db.query('ALTER TABLE syncs DROP COLUMN next_due_at').run();
- }
- }).immediate();
+ runMigrations(db);
return db;
} catch (error) {
db.close();
diff --git a/packages/sync/src/db/migrate.ts b/packages/sync/src/db/migrate.ts
new file mode 100644
index 0000000..7363a8b
--- /dev/null
+++ b/packages/sync/src/db/migrate.ts
@@ -0,0 +1,32 @@
+import type { Database } from 'bun:sqlite';
+import { up as initialSchema } from './migrations/0000-initial-schema';
+import { up as centralizePolling } from './migrations/0001-centralize-polling';
+
+// Static imports keep migration history available in both installed packages and compiled hosts.
+const migrations = [
+ { name: '0000-initial-schema', up: initialSchema },
+ { name: '0001-centralize-polling', up: centralizePolling },
+] as const;
+
+export function runMigrations(db: Database): void {
+ db.transaction(() => {
+ db.query(`CREATE TABLE IF NOT EXISTS __migrations (
+ name TEXT PRIMARY KEY, applied_at TEXT NOT NULL
+ )`).run();
+ const applied = db
+ .query<{ name: string }, []>('SELECT name FROM __migrations ORDER BY name')
+ .all();
+ for (const [index, { name }] of applied.entries()) {
+ if (name !== migrations[index]?.name) {
+ throw new Error('Unsupported sync database migration history');
+ }
+ }
+ for (const migration of migrations.slice(applied.length)) {
+ migration.up(db);
+ db.query('INSERT INTO __migrations (name,applied_at) VALUES (?,?)').run(
+ migration.name,
+ new Date().toISOString(),
+ );
+ }
+ }).immediate();
+}
diff --git a/packages/sync/src/db/migrations/0000-initial-schema.ts b/packages/sync/src/db/migrations/0000-initial-schema.ts
new file mode 100644
index 0000000..fefcf87
--- /dev/null
+++ b/packages/sync/src/db/migrations/0000-initial-schema.ts
@@ -0,0 +1,10 @@
+import type { Database } from 'bun:sqlite';
+import schema from '../schema.sql' with { type: 'text' };
+
+export function up(db: Database): void {
+ // The released schema is idempotent, so existing unversioned databases adopt this baseline.
+ // Its simple DDL statements end with semicolon/newline; execute separately to catch each failure.
+ for (const statement of schema.split(';\n').filter((sql) => sql.trim())) {
+ db.query(statement).run();
+ }
+}
diff --git a/packages/sync/src/db/migrations/0001-centralize-polling.ts b/packages/sync/src/db/migrations/0001-centralize-polling.ts
new file mode 100644
index 0000000..04a8e88
--- /dev/null
+++ b/packages/sync/src/db/migrations/0001-centralize-polling.ts
@@ -0,0 +1,9 @@
+import type { Database } from 'bun:sqlite';
+
+export function up(db: Database): void {
+ db.query('ALTER TABLE syncs ADD COLUMN retry_at INTEGER').run();
+ db.query(`UPDATE syncs SET retry_at=next_due_at
+ WHERE status IN ('retrying','interrupted','waiting_for_capacity')`).run();
+ db.query('ALTER TABLE syncs DROP COLUMN interval_ms').run();
+ db.query('ALTER TABLE syncs DROP COLUMN next_due_at').run();
+}
diff --git a/packages/sync/src/db/schema.sql b/packages/sync/src/db/schema.sql
index 469ce53..7a288f1 100644
--- a/packages/sync/src/db/schema.sql
+++ b/packages/sync/src/db/schema.sql
@@ -3,7 +3,7 @@ CREATE TABLE IF NOT EXISTS syncs (
connection TEXT, config TEXT NOT NULL, destination_type TEXT NOT NULL, destination_config TEXT NOT NULL,
enabled INTEGER NOT NULL,
checkpoint TEXT NOT NULL,
- retry_at INTEGER, status TEXT NOT NULL,
+ interval_ms INTEGER NOT NULL, next_due_at INTEGER NOT NULL, status TEXT NOT NULL,
error_code TEXT,
generation INTEGER NOT NULL DEFAULT 0, expires_at INTEGER,
failure_count INTEGER NOT NULL DEFAULT 0, resync INTEGER NOT NULL DEFAULT 0,
diff --git a/packages/sync/test/polling-migration.test.ts b/packages/sync/test/polling-migration.test.ts
index 3838842..01326f8 100644
--- a/packages/sync/test/polling-migration.test.ts
+++ b/packages/sync/test/polling-migration.test.ts
@@ -1,26 +1,49 @@
import { Database } from 'bun:sqlite';
import { expect, test } from 'bun:test';
import { openDatabase } from '../src/db/client';
+import releasedSchema from '../src/db/schema.sql' with { type: 'text' };
+import type { Deliverable } from '../src/models/delivery';
import { defaultLimits } from '../src/models/limits';
import { SqliteAcquisition } from '../src/repositories/acquisition/sqlite';
import { SqliteCatalog } from '../src/repositories/catalog/sqlite';
+import { SqliteDeliveries } from '../src/repositories/delivery/sqlite';
import { alpha, fixture, storage } from './support';
-const legacySchema = `CREATE TABLE syncs (
- owner_id TEXT NOT NULL, id TEXT NOT NULL, definition_id TEXT NOT NULL,
- connection TEXT, config TEXT NOT NULL, destination_type TEXT NOT NULL, destination_config TEXT NOT NULL,
- enabled INTEGER NOT NULL, checkpoint TEXT NOT NULL,
- interval_ms INTEGER NOT NULL, next_due_at INTEGER NOT NULL, status TEXT NOT NULL,
- error_code TEXT, generation INTEGER NOT NULL DEFAULT 0, expires_at INTEGER,
- failure_count INTEGER NOT NULL DEFAULT 0, resync INTEGER NOT NULL DEFAULT 0,
- PRIMARY KEY(owner_id,id)
-);`;
+const migrationNames = ['0000-initial-schema', '0001-centralize-polling'];
const leaseMs = 60_000;
const futureRetry = 9_000_000_000_000;
+const retryStates = ['retrying', 'interrupted', 'waiting_for_capacity'];
+const delivery: Omit = {
+ id: 'saved-delivery',
+ ownerId: alpha.ownerId,
+ syncId: 'succeeded',
+ definition: fixture.definition.id,
+ records: [
+ {
+ operation: 'upsert',
+ kind: 'item',
+ id: 'saved',
+ revision: 1,
+ data: { value: 1 },
+ assetRefs: { file: { id: 'saved-asset', version: '1' } },
+ },
+ ],
+ assets: [
+ {
+ id: 'saved-asset',
+ version: '1',
+ name: 'saved.txt',
+ mediaType: 'text/plain',
+ size: Buffer.byteLength('saved'),
+ sha256: new Bun.CryptoHasher('sha256').update('saved').digest('hex'),
+ },
+ ],
+};
function legacyDatabase(path: string) {
const db = new Database(path);
- db.exec(legacySchema);
- for (const state of ['succeeded', 'retrying', 'disabled', 'ready', 'running']) {
+ db.exec(releasedSchema);
+ db.exec('PRAGMA foreign_keys=ON');
+ for (const state of ['succeeded', 'disabled', 'ready', 'running', ...retryStates]) {
db.query(`INSERT INTO syncs (owner_id,id,definition_id,config,destination_type,
destination_config,enabled,checkpoint,interval_ms,next_due_at,status,generation,expires_at,
failure_count,resync) VALUES (?,?,?,'{"count":3}','local','{}',?, '2',60000,?,?,7,?,2,1)`).run(
@@ -33,18 +56,48 @@ function legacyDatabase(path: string) {
state === 'running' ? futureRetry : null,
);
}
+ db.query(
+ `INSERT INTO record_state VALUES ('alpha','succeeded','item','saved','saved-hash',1,0)`,
+ ).run();
+ db.query(`INSERT INTO sync_polls(owner_id,sync_id,started_at,completed_at,records_processed,records_queued,state)
+ VALUES ('alpha','succeeded',1,2,1,1,'succeeded')`).run();
+ db.query(`INSERT INTO sync_polls(owner_id,sync_id,started_at,state,error_code)
+ VALUES ('alpha','retrying',1,'retrying','execution_failed')`).run();
+ const body = JSON.stringify(delivery);
+ db.query(`INSERT INTO deliveries(owner_id,id,sync_id,body,bytes,record_count,due_at)
+ VALUES (?,?,?,?,?,1,0)`).run(
+ alpha.ownerId,
+ delivery.id,
+ delivery.syncId,
+ body,
+ Buffer.byteLength(body),
+ );
+ db.query(`INSERT INTO delivery_assets(id,owner_id,sync_id,generation,asset_id,asset_version,
+ descriptor,ready,bytes,delivery_id) VALUES ('queued-file','alpha','succeeded',7,'saved-asset','1',?,1,5,?)`).run(
+ JSON.stringify(delivery.assets[0]),
+ delivery.id,
+ );
return db;
}
+function queuedState(db: Database) {
+ return {
+ records: db.query('SELECT * FROM record_state').all(),
+ polls: db.query('SELECT * FROM sync_polls ORDER BY id').all(),
+ deliveries: db.query('SELECT * FROM deliveries ORDER BY sequence').all(),
+ assets: db.query('SELECT * FROM delivery_assets').all(),
+ };
+}
+function history(db: Database) {
+ return db
+ .query<{ name: string; applied_at: string }, []>('SELECT * FROM __migrations ORDER BY name')
+ .all();
+}
-test('legacy polling schedules migrate without losing progress, leases, backoff or queued data', () => {
+test('released databases adopt migration history without losing progress, leases, retries or queued data', () => {
const files = storage();
let db = legacyDatabase(files.path);
- // Real foreign keys and populated dependent tables must survive the column removal.
- db.exec(`CREATE TABLE retained (
- sync_id TEXT, owner_id TEXT, payload TEXT,
- FOREIGN KEY(owner_id,sync_id) REFERENCES syncs(owner_id,id)
- ); INSERT INTO retained VALUES ('succeeded','alpha','queued records and assets');`);
const before = db.query('SELECT * FROM syncs ORDER BY id').all() as Record[];
+ const queue = queuedState(db);
db.close();
try {
db = openDatabase(files.path);
@@ -52,16 +105,19 @@ test('legacy polling schedules migrate without losing progress, leases, backoff
expect(after).toEqual(
before.map(({ interval_ms: _interval, next_due_at, ...row }) => ({
...row,
- retry_at: row.status === 'retrying' ? next_due_at : null,
+ retry_at: retryStates.includes(String(row.status)) ? next_due_at : null,
})),
);
expect(db.query('PRAGMA foreign_key_check').all()).toEqual([]);
- expect(db.query('SELECT * FROM retained').all()).toEqual([
- { sync_id: 'succeeded', owner_id: 'alpha', payload: 'queued records and assets' },
- ]);
+ expect(queuedState(db)).toEqual(queue);
+ const applied = history(db);
+ expect(applied.map(({ name }) => name)).toEqual(migrationNames);
+ // Reopening must leave both the data and recorded application times unchanged.
db.close();
db = openDatabase(files.path);
+ expect(history(db)).toEqual(applied);
expect(db.query('SELECT * FROM syncs ORDER BY id').all()).toEqual(after);
+ expect(queuedState(db)).toEqual(queue);
const catalog = new SqliteCatalog(db);
const acquisition = new SqliteAcquisition({ db, limits: defaultLimits });
expect(acquisition.claim(leaseMs)?.sync.id).toBe('ready');
@@ -69,6 +125,7 @@ test('legacy polling schedules migrate without losing progress, leases, backoff
acquisition.poll();
expect(acquisition.claim(leaseMs)?.sync.id).toBe('succeeded');
expect(acquisition.claim(leaseMs)).toBeUndefined();
+ expect(new SqliteDeliveries(db).claim(leaseMs)?.delivery).toEqual(delivery);
expect(
catalog.createSync({
...alpha,
@@ -84,25 +141,99 @@ test('legacy polling schedules migrate without losing progress, leases, backoff
}
});
-test('failed migration rolls back all schema and data changes and can be retried', () => {
+test('fresh databases apply the same ordered migration history', () => {
+ using db = openDatabase(':memory:');
+ expect(history(db).map(({ name }) => name)).toEqual(migrationNames);
+ const columns = db
+ .query<{ name: string }, []>("SELECT name FROM pragma_table_info('syncs')")
+ .all()
+ .map(({ name }) => name);
+ expect(columns).toContain('retry_at');
+ expect(columns).not.toContain('interval_ms');
+ expect(columns).not.toContain('next_due_at');
+});
+
+test.each([false, true])(
+ 'failed migration rolls back schema, data and history (baseline recorded: %s)',
+ (recorded) => {
+ const files = storage();
+ let db = legacyDatabase(files.path);
+ if (recorded) {
+ db.query('CREATE TABLE __migrations (name TEXT PRIMARY KEY, applied_at TEXT NOT NULL)').run();
+ db.query('INSERT INTO __migrations VALUES (?,?)').run(
+ migrationNames[0]!,
+ '2026-10-01T00:00:00.000Z',
+ );
+ }
+ db.exec(
+ "CREATE TRIGGER fail_migration BEFORE UPDATE ON syncs BEGIN SELECT RAISE(ABORT,'injected'); END",
+ );
+ const before = db.query('SELECT * FROM syncs ORDER BY id').all();
+ const schemaBefore = db.query('SELECT * FROM sqlite_schema ORDER BY name').all();
+ const queue = queuedState(db);
+ const historyBefore = recorded ? history(db) : [];
+ db.close();
+ try {
+ expect(() => openDatabase(files.path)).toThrow('injected');
+ db = new Database(files.path);
+ expect(db.query('SELECT * FROM syncs ORDER BY id').all()).toEqual(before);
+ expect(db.query('SELECT * FROM sqlite_schema ORDER BY name').all()).toEqual(schemaBefore);
+ expect(queuedState(db)).toEqual(queue);
+ if (recorded) {
+ expect(history(db)).toEqual(historyBefore);
+ }
+ db.exec('DROP TRIGGER fail_migration');
+ db.close();
+ db = openDatabase(files.path);
+ expect(history(db).map(({ name }) => name)).toEqual(migrationNames);
+ expect(db.query('SELECT id FROM syncs').all()).toHaveLength(before.length);
+ expect(queuedState(db)).toEqual(queue);
+ } finally {
+ db.close();
+ files.close();
+ }
+ },
+);
+
+test.each(['future', 'missing-baseline'])(
+ 'unsupported migration history fails closed (%s)',
+ (state) => {
+ const files = storage();
+ let db = openDatabase(files.path);
+ if (state === 'future') {
+ db.query('INSERT INTO __migrations VALUES (?,?)').run(
+ '9999_future',
+ new Date().toISOString(),
+ );
+ } else {
+ db.query('DELETE FROM __migrations WHERE name=?').run(migrationNames[0]!);
+ }
+ const before = history(db);
+ const schemaBefore = db.query('SELECT * FROM sqlite_schema ORDER BY name').all();
+ db.close();
+ try {
+ expect(() => openDatabase(files.path)).toThrow('Unsupported sync database migration history');
+ db = new Database(files.path);
+ expect(history(db)).toEqual(before);
+ expect(db.query('SELECT * FROM sqlite_schema ORDER BY name').all()).toEqual(schemaBefore);
+ } finally {
+ db.close();
+ files.close();
+ }
+ },
+);
+
+test('an incompatible unversioned database fails without recording a partial baseline', () => {
const files = storage();
- let db = legacyDatabase(files.path);
- db.exec(
- "CREATE TRIGGER fail_migration BEFORE UPDATE ON syncs BEGIN SELECT RAISE(ABORT,'injected'); END",
- );
- const before = db.query('SELECT * FROM syncs ORDER BY id').all();
+ let db = new Database(files.path);
+ // A failure in the middle of the baseline must not be hidden by later successful DDL.
+ db.query('CREATE TABLE deliveries (unexpected TEXT)').run();
+ const before = db.query('SELECT * FROM sqlite_schema ORDER BY name').all();
db.close();
try {
expect(() => openDatabase(files.path)).toThrow();
db = new Database(files.path);
- expect(db.query('SELECT * FROM syncs ORDER BY id').all()).toEqual(before);
- expect(db.query("SELECT name FROM pragma_table_info('syncs')").all()).not.toContainEqual({
- name: 'retry_at',
- });
- db.exec('DROP TRIGGER fail_migration');
- db.close();
- db = openDatabase(files.path);
- expect(db.query('SELECT id FROM syncs').all()).toHaveLength(before.length);
+ expect(db.query('SELECT * FROM sqlite_schema ORDER BY name').all()).toEqual(before);
} finally {
db.close();
files.close();
From 0c84aa84a666735e9af8a36444d60d21403c841c Mon Sep 17 00:00:00 2001
From: massimoalbarello
Date: Fri, 2 Oct 2026 19:44:54 +0100
Subject: [PATCH 3/5] refactor: use SQL migrations and retry sources on cron
---
packages/sync/src/db/migrate.ts | 14 +-
.../src/db/migrations/0000-initial-schema.ts | 10 -
.../0000-initial.sql} | 10 +
.../0000-initial.sql.d.ts} | 0
.../db/migrations/0001-centralize-polling.sql | 5 +
.../0001-centralize-polling.sql.d.ts | 2 +
.../db/migrations/0001-centralize-polling.ts | 9 -
packages/sync/src/models/sync.ts | 3 -
.../src/repositories/acquisition/contract.ts | 3 -
.../src/repositories/acquisition/lease.ts | 20 +-
.../src/repositories/acquisition/sqlite.ts | 8 +-
.../sync/src/repositories/assets/sqlite.ts | 2 +-
.../sync/src/repositories/catalog/sqlite.ts | 12 +-
.../sync/src/repositories/delivery/sqlite.ts | 2 +-
packages/sync/src/repositories/rows.ts | 1 -
packages/sync/src/services/acquisition.ts | 18 +-
packages/sync/test/acquisition-retry.test.ts | 186 ++++--------------
packages/sync/test/capacity.test.ts | 2 +-
packages/sync/test/connector-failure.test.ts | 22 +--
...ckoff.test.ts => delivery-backoff.test.ts} | 16 +-
packages/sync/test/persistence.test.ts | 2 +-
packages/sync/test/polling-migration.test.ts | 35 ++--
packages/sync/test/run-due.test.ts | 18 +-
packages/sync/test/source-http-error.test.ts | 33 ++--
24 files changed, 142 insertions(+), 291 deletions(-)
delete mode 100644 packages/sync/src/db/migrations/0000-initial-schema.ts
rename packages/sync/src/db/{schema.sql => migrations/0000-initial.sql} (92%)
rename packages/sync/src/db/{schema.sql.d.ts => migrations/0000-initial.sql.d.ts} (100%)
create mode 100644 packages/sync/src/db/migrations/0001-centralize-polling.sql
create mode 100644 packages/sync/src/db/migrations/0001-centralize-polling.sql.d.ts
delete mode 100644 packages/sync/src/db/migrations/0001-centralize-polling.ts
rename packages/sync/test/{shared-backoff.test.ts => delivery-backoff.test.ts} (69%)
diff --git a/packages/sync/src/db/migrate.ts b/packages/sync/src/db/migrate.ts
index 7363a8b..cf15efc 100644
--- a/packages/sync/src/db/migrate.ts
+++ b/packages/sync/src/db/migrate.ts
@@ -1,11 +1,11 @@
import type { Database } from 'bun:sqlite';
-import { up as initialSchema } from './migrations/0000-initial-schema';
-import { up as centralizePolling } from './migrations/0001-centralize-polling';
+import initialSchema from './migrations/0000-initial.sql' with { type: 'text' };
+import centralizePolling from './migrations/0001-centralize-polling.sql' with { type: 'text' };
// Static imports keep migration history available in both installed packages and compiled hosts.
const migrations = [
- { name: '0000-initial-schema', up: initialSchema },
- { name: '0001-centralize-polling', up: centralizePolling },
+ { name: '0000-initial.sql', sql: initialSchema },
+ { name: '0001-centralize-polling.sql', sql: centralizePolling },
] as const;
export function runMigrations(db: Database): void {
@@ -22,7 +22,11 @@ export function runMigrations(db: Database): void {
}
}
for (const migration of migrations.slice(applied.length)) {
- migration.up(db);
+ // Explicit boundaries also support triggers and strings containing semicolons.
+ // Executing each statement separately surfaces errors hidden by Bun 1.4 SQL batches.
+ for (const statement of migration.sql.split('--> statement-breakpoint')) {
+ db.query(statement).run();
+ }
db.query('INSERT INTO __migrations (name,applied_at) VALUES (?,?)').run(
migration.name,
new Date().toISOString(),
diff --git a/packages/sync/src/db/migrations/0000-initial-schema.ts b/packages/sync/src/db/migrations/0000-initial-schema.ts
deleted file mode 100644
index fefcf87..0000000
--- a/packages/sync/src/db/migrations/0000-initial-schema.ts
+++ /dev/null
@@ -1,10 +0,0 @@
-import type { Database } from 'bun:sqlite';
-import schema from '../schema.sql' with { type: 'text' };
-
-export function up(db: Database): void {
- // The released schema is idempotent, so existing unversioned databases adopt this baseline.
- // Its simple DDL statements end with semicolon/newline; execute separately to catch each failure.
- for (const statement of schema.split(';\n').filter((sql) => sql.trim())) {
- db.query(statement).run();
- }
-}
diff --git a/packages/sync/src/db/schema.sql b/packages/sync/src/db/migrations/0000-initial.sql
similarity index 92%
rename from packages/sync/src/db/schema.sql
rename to packages/sync/src/db/migrations/0000-initial.sql
index 7a288f1..2dbd1b3 100644
--- a/packages/sync/src/db/schema.sql
+++ b/packages/sync/src/db/migrations/0000-initial.sql
@@ -9,12 +9,14 @@ CREATE TABLE IF NOT EXISTS syncs (
failure_count INTEGER NOT NULL DEFAULT 0, resync INTEGER NOT NULL DEFAULT 0,
PRIMARY KEY(owner_id, id)
);
+--> statement-breakpoint
CREATE TABLE IF NOT EXISTS record_state (
owner_id TEXT NOT NULL, sync_id TEXT NOT NULL, kind TEXT NOT NULL, id TEXT NOT NULL,
hash TEXT NOT NULL, revision INTEGER NOT NULL, deleted INTEGER NOT NULL,
PRIMARY KEY(owner_id, sync_id, kind, id),
FOREIGN KEY(owner_id, sync_id) REFERENCES syncs(owner_id, id)
);
+--> statement-breakpoint
-- Polling diagnostics only; execution leases and scheduling belong to syncs.
CREATE TABLE IF NOT EXISTS sync_polls (
id INTEGER PRIMARY KEY AUTOINCREMENT, owner_id TEXT NOT NULL, sync_id TEXT NOT NULL,
@@ -23,8 +25,11 @@ CREATE TABLE IF NOT EXISTS sync_polls (
state TEXT NOT NULL, error_code TEXT,
FOREIGN KEY(owner_id, sync_id) REFERENCES syncs(owner_id, id) ON DELETE CASCADE
);
+--> statement-breakpoint
CREATE INDEX IF NOT EXISTS sync_poll_history ON sync_polls(owner_id,sync_id,id);
+--> statement-breakpoint
CREATE UNIQUE INDEX IF NOT EXISTS sync_poll_active ON sync_polls(owner_id,sync_id) WHERE completed_at IS NULL;
+--> statement-breakpoint
CREATE TABLE IF NOT EXISTS deliveries (
sequence INTEGER PRIMARY KEY AUTOINCREMENT, owner_id TEXT NOT NULL, id TEXT NOT NULL,
sync_id TEXT NOT NULL, body TEXT NOT NULL,
@@ -34,8 +39,11 @@ CREATE TABLE IF NOT EXISTS deliveries (
UNIQUE(owner_id, id),
FOREIGN KEY(owner_id, sync_id) REFERENCES syncs(owner_id, id)
);
+--> statement-breakpoint
CREATE INDEX IF NOT EXISTS delivery_order ON deliveries(owner_id,sync_id,sequence);
+--> statement-breakpoint
CREATE INDEX IF NOT EXISTS deliveries_due ON deliveries(state, due_at);
+--> statement-breakpoint
-- Staging, queued bytes, and pending cleanup share one delivery-owned ledger.
-- No foreign keys: cleanup must survive deletion of the delivery or sync.
CREATE TABLE IF NOT EXISTS delivery_assets (
@@ -46,5 +54,7 @@ CREATE TABLE IF NOT EXISTS delivery_assets (
bytes INTEGER NOT NULL DEFAULT 0 CHECK(bytes>=0), delivery_id TEXT,
UNIQUE(owner_id,sync_id,generation,asset_id,asset_version)
);
+--> statement-breakpoint
CREATE INDEX IF NOT EXISTS delivery_asset_queue ON delivery_assets(owner_id,delivery_id);
+--> statement-breakpoint
CREATE INDEX IF NOT EXISTS delivery_asset_sync ON delivery_assets(owner_id,sync_id);
diff --git a/packages/sync/src/db/schema.sql.d.ts b/packages/sync/src/db/migrations/0000-initial.sql.d.ts
similarity index 100%
rename from packages/sync/src/db/schema.sql.d.ts
rename to packages/sync/src/db/migrations/0000-initial.sql.d.ts
diff --git a/packages/sync/src/db/migrations/0001-centralize-polling.sql b/packages/sync/src/db/migrations/0001-centralize-polling.sql
new file mode 100644
index 0000000..2eed189
--- /dev/null
+++ b/packages/sync/src/db/migrations/0001-centralize-polling.sql
@@ -0,0 +1,5 @@
+ALTER TABLE syncs DROP COLUMN interval_ms;
+--> statement-breakpoint
+ALTER TABLE syncs DROP COLUMN next_due_at;
+--> statement-breakpoint
+ALTER TABLE syncs DROP COLUMN failure_count;
diff --git a/packages/sync/src/db/migrations/0001-centralize-polling.sql.d.ts b/packages/sync/src/db/migrations/0001-centralize-polling.sql.d.ts
new file mode 100644
index 0000000..6fe9835
--- /dev/null
+++ b/packages/sync/src/db/migrations/0001-centralize-polling.sql.d.ts
@@ -0,0 +1,2 @@
+declare const sql: string;
+export default sql;
diff --git a/packages/sync/src/db/migrations/0001-centralize-polling.ts b/packages/sync/src/db/migrations/0001-centralize-polling.ts
deleted file mode 100644
index 04a8e88..0000000
--- a/packages/sync/src/db/migrations/0001-centralize-polling.ts
+++ /dev/null
@@ -1,9 +0,0 @@
-import type { Database } from 'bun:sqlite';
-
-export function up(db: Database): void {
- db.query('ALTER TABLE syncs ADD COLUMN retry_at INTEGER').run();
- db.query(`UPDATE syncs SET retry_at=next_due_at
- WHERE status IN ('retrying','interrupted','waiting_for_capacity')`).run();
- db.query('ALTER TABLE syncs DROP COLUMN interval_ms').run();
- db.query('ALTER TABLE syncs DROP COLUMN next_due_at').run();
-}
diff --git a/packages/sync/src/models/sync.ts b/packages/sync/src/models/sync.ts
index 018b6da..f05d120 100644
--- a/packages/sync/src/models/sync.ts
+++ b/packages/sync/src/models/sync.ts
@@ -11,7 +11,6 @@ export interface Sync {
destination: { type: string; config: JsonObject };
enabled: boolean;
checkpoint: JsonValue;
- retryAt: number | null;
status: SyncStatus;
errorCode: string | null;
}
@@ -49,7 +48,6 @@ export interface SyncSummary {
destinationType: string;
connection?: ConnectionRef;
enabled: boolean;
- retryAt: number | null;
status: SyncStatus;
errorCode: string | null;
}
@@ -60,7 +58,6 @@ export function summarizeSync(sync: Sync): SyncSummary {
destinationType: sync.destination.type,
connection: sync.connection,
enabled: sync.enabled,
- retryAt: sync.retryAt,
status: sync.status,
errorCode: sync.errorCode,
};
diff --git a/packages/sync/src/repositories/acquisition/contract.ts b/packages/sync/src/repositories/acquisition/contract.ts
index e65952e..be280c3 100644
--- a/packages/sync/src/repositories/acquisition/contract.ts
+++ b/packages/sync/src/repositories/acquisition/contract.ts
@@ -6,7 +6,6 @@ export interface AcquisitionLease extends Scope {
sync: Sync;
generation: number;
force: boolean;
- failureCount: number;
}
export interface AcquisitionRepository {
poll(): void;
@@ -18,8 +17,6 @@ export interface AcquisitionRepository {
lease: AcquisitionLease;
state: Exclude;
errorCode?: string;
- delay: number;
- failureCount?: number;
pause?: boolean;
}): void;
}
diff --git a/packages/sync/src/repositories/acquisition/lease.ts b/packages/sync/src/repositories/acquisition/lease.ts
index a7c230b..1c9ef4a 100644
--- a/packages/sync/src/repositories/acquisition/lease.ts
+++ b/packages/sync/src/repositories/acquisition/lease.ts
@@ -30,7 +30,7 @@ export function claimAcquisition({
WHERE completed_at IS NULL AND EXISTS (SELECT 1 FROM syncs
WHERE syncs.owner_id=sync_polls.owner_id AND syncs.id=sync_polls.sync_id
AND enabled=1 AND expires_at<=?)`).run(now);
- db.query(`UPDATE syncs SET retry_at=NULL,status='interrupted',error_code='lease_expired',expires_at=NULL
+ db.query(`UPDATE syncs SET status='ready',error_code='lease_expired',expires_at=NULL
WHERE enabled=1 AND expires_at<=?`).run(now);
const row = db
.query<
@@ -39,13 +39,12 @@ export function claimAcquisition({
owner_id: string;
generation: number;
resync: number;
- failure_count: number;
},
- [number]
- >(`SELECT id,owner_id,generation,resync,failure_count FROM syncs
- WHERE enabled=1 AND status!='succeeded' AND expires_at IS NULL AND COALESCE(retry_at,0)<=?
- ORDER BY COALESCE(retry_at,0),generation,id LIMIT 1`)
- .get(now);
+ []
+ >(`SELECT id,owner_id,generation,resync FROM syncs
+ WHERE enabled=1 AND status='ready' AND expires_at IS NULL
+ ORDER BY generation,id LIMIT 1`)
+ .get();
if (!row) {
return;
}
@@ -69,7 +68,6 @@ export function claimAcquisition({
sync: readSync({ db, scope: { ...scope, id: row.id } }),
generation: row.generation + 1,
force: row.resync === 1,
- failureCount: row.failure_count,
};
})
.immediate();
@@ -88,15 +86,13 @@ export function finishAcquisition(
lease.ownerId,
lease.sync.id,
);
- db.query(`UPDATE syncs SET status=?,error_code=?,retry_at=?,expires_at=NULL,
+ db.query(`UPDATE syncs SET status=?,error_code=?,expires_at=NULL,
enabled=CASE WHEN ? THEN 0 ELSE enabled END,
- failure_count=COALESCE(?,failure_count),resync=CASE WHEN ? THEN 0 ELSE resync END
+ resync=CASE WHEN ? THEN 0 ELSE resync END
WHERE owner_id=? AND id=?`).run(
input.pause ? 'disabled' : input.state,
input.errorCode ?? null,
- input.delay > 0 ? now + input.delay : null,
Number(input.pause ?? false),
- input.failureCount ?? null,
Number(input.state === 'succeeded'),
lease.ownerId,
lease.sync.id,
diff --git a/packages/sync/src/repositories/acquisition/sqlite.ts b/packages/sync/src/repositories/acquisition/sqlite.ts
index 7a72447..230d0e2 100644
--- a/packages/sync/src/repositories/acquisition/sqlite.ts
+++ b/packages/sync/src/repositories/acquisition/sqlite.ts
@@ -17,12 +17,14 @@ export class SqliteAcquisition implements AcquisitionRepository {
constructor(private readonly input: { db: Database; limits: QueueLimits }) {}
poll(): void {
this.input.db
- .query("UPDATE syncs SET status='ready' WHERE enabled=1 AND status='succeeded'")
+ .query(
+ "UPDATE syncs SET status='ready' WHERE enabled=1 AND status IN ('succeeded','retrying','interrupted','waiting_for_capacity')",
+ )
.run();
}
capacityReleased(): void {
this.input.db
- .query("UPDATE syncs SET retry_at=NULL WHERE enabled=1 AND status='waiting_for_capacity'")
+ .query("UPDATE syncs SET status='ready' WHERE enabled=1 AND status='waiting_for_capacity'")
.run();
}
claim(leaseMs: number) {
@@ -57,8 +59,6 @@ export class SqliteAcquisition implements AcquisitionRepository {
db,
lease: input.lease,
state: page.complete ? 'succeeded' : 'ready',
- delay: 0,
- failureCount: page.complete ? 0 : undefined,
});
}).immediate();
}
diff --git a/packages/sync/src/repositories/assets/sqlite.ts b/packages/sync/src/repositories/assets/sqlite.ts
index f4d15c9..01d50fb 100644
--- a/packages/sync/src/repositories/assets/sqlite.ts
+++ b/packages/sync/src/repositories/assets/sqlite.ts
@@ -103,7 +103,7 @@ export class SqliteAssets implements AssetRepository {
if (deleted) {
this.input.db
.query(
- "UPDATE syncs SET retry_at=NULL WHERE enabled=1 AND status='waiting_for_capacity' AND NOT (owner_id=? AND id=?)",
+ "UPDATE syncs SET status='ready' WHERE enabled=1 AND status='waiting_for_capacity' AND NOT (owner_id=? AND id=?)",
)
.run(deleted.owner_id, deleted.sync_id);
}
diff --git a/packages/sync/src/repositories/catalog/sqlite.ts b/packages/sync/src/repositories/catalog/sqlite.ts
index 8ab2280..eb474c5 100644
--- a/packages/sync/src/repositories/catalog/sqlite.ts
+++ b/packages/sync/src/repositories/catalog/sqlite.ts
@@ -65,7 +65,7 @@ export class SqliteCatalog implements CatalogRepository {
}
this.db
.query(`UPDATE syncs SET connection=?,enabled=1,
- status='ready',error_code=NULL,retry_at=NULL WHERE owner_id=? AND id=?`)
+ status='ready',error_code=NULL WHERE owner_id=? AND id=?`)
.run(canonicalJson(input.connection).json, input.ownerId, input.id);
return this.sync(input);
})
@@ -84,7 +84,7 @@ export class SqliteCatalog implements CatalogRepository {
.run(input.enabled ? 'ready' : 'disabled', input.ownerId, input.id);
this.db
.query(
- `UPDATE syncs SET enabled=?,status=?,error_code=NULL,retry_at=NULL,generation=generation+1,expires_at=NULL WHERE owner_id=? AND id=?`,
+ `UPDATE syncs SET enabled=?,status=?,error_code=NULL,generation=generation+1,expires_at=NULL WHERE owner_id=? AND id=?`,
)
.run(
Number(input.enabled),
@@ -110,7 +110,7 @@ export class SqliteCatalog implements CatalogRepository {
fail('busy');
}
this.db
- .query(`UPDATE syncs SET status='ready',retry_at=NULL WHERE owner_id=? AND id=?`)
+ .query(`UPDATE syncs SET status='ready' WHERE owner_id=? AND id=?`)
.run(input.ownerId, input.id);
})
.immediate();
@@ -127,7 +127,7 @@ export class SqliteCatalog implements CatalogRepository {
.run(Date.now(), input.ownerId, input.id);
this.db
.query(
- `UPDATE syncs SET checkpoint=?,status='ready',error_code=NULL,retry_at=NULL,resync=1,failure_count=0,generation=generation+1,expires_at=NULL WHERE owner_id=? AND id=?`,
+ `UPDATE syncs SET checkpoint=?,status='ready',error_code=NULL,resync=1,generation=generation+1,expires_at=NULL WHERE owner_id=? AND id=?`,
)
.run(canonicalJson(input.checkpoint).json, input.ownerId, input.id);
})
@@ -145,7 +145,9 @@ export class SqliteCatalog implements CatalogRepository {
.run(input.ownerId, input.id);
this.db.query('DELETE FROM syncs WHERE owner_id=? AND id=?').run(input.ownerId, input.id);
this.db
- .query(`UPDATE syncs SET retry_at=NULL WHERE enabled=1 AND status='waiting_for_capacity'`)
+ .query(
+ `UPDATE syncs SET status='ready' WHERE enabled=1 AND status='waiting_for_capacity'`,
+ )
.run();
})
.immediate();
diff --git a/packages/sync/src/repositories/delivery/sqlite.ts b/packages/sync/src/repositories/delivery/sqlite.ts
index 798f90b..d93de4e 100644
--- a/packages/sync/src/repositories/delivery/sqlite.ts
+++ b/packages/sync/src/repositories/delivery/sqlite.ts
@@ -59,7 +59,7 @@ export class SqliteDeliveries implements DeliveryRepository {
// Capacity is global; all paused acquisitions can compete for the released budget.
this.db
.query(
- "UPDATE syncs SET retry_at=NULL WHERE enabled=1 AND status='waiting_for_capacity'",
+ "UPDATE syncs SET status='ready' WHERE enabled=1 AND status='waiting_for_capacity'",
)
.run();
} else {
diff --git a/packages/sync/src/repositories/rows.ts b/packages/sync/src/repositories/rows.ts
index 0eb89c1..6578b08 100644
--- a/packages/sync/src/repositories/rows.ts
+++ b/packages/sync/src/repositories/rows.ts
@@ -23,7 +23,6 @@ export function readSync(input: { db: Database; scope: Resource }): Sync {
},
enabled: row.enabled === 1,
checkpoint: JSON.parse(String(row.checkpoint)),
- retryAt: row.retry_at === null ? null : Number(row.retry_at),
status: row.status as SyncStatus,
errorCode: row.error_code === null ? null : String(row.error_code),
};
diff --git a/packages/sync/src/services/acquisition.ts b/packages/sync/src/services/acquisition.ts
index f97c78c..4cef3d6 100644
--- a/packages/sync/src/services/acquisition.ts
+++ b/packages/sync/src/services/acquisition.ts
@@ -4,7 +4,7 @@ import { bindProvider, type ProviderGateway } from '../execution/provider';
import { assetKey, type SourceAssets } from '../models/asset';
import { SyncError } from '../models/error';
import { canonicalJson } from '../models/json';
-import { retryDelay, type Timing } from '../models/limits';
+import type { Timing } from '../models/limits';
import type { Registry } from '../models/registry';
import { retryableStatus } from '../models/source-http-error';
import { identifier } from '../models/validation';
@@ -34,11 +34,11 @@ export class AcquisitionService {
return this.input.repository.claim(this.input.timing.leaseMs);
}
async execute(input: { lease: AcquisitionLease; signal: AbortSignal }): Promise {
- const { repository, timing } = this.input;
+ const { repository } = this.input;
const { lease } = input;
try {
if (!repository.hasCapacity()) {
- this.finish({ lease, state: 'waiting_for_capacity', delay: timing.retryMs });
+ this.finish({ lease, state: 'waiting_for_capacity' });
return;
}
await abortable({ signal: input.signal, run: () => this.consume(input) });
@@ -47,20 +47,15 @@ export class AcquisitionService {
}
}
private failed(input: { lease: AcquisitionLease; signal: AbortSignal; error: unknown }): void {
- const { timing, log } = this.input;
+ const { log } = this.input;
const { lease, signal, error } = input;
const code = signal.aborted
? abortCode(signal)
: error instanceof SyncError
? error.code
: 'execution_failed';
- const failed = !['paused', 'interrupted', 'waiting_for_capacity'].includes(code);
- const failureCount = failed ? lease.failureCount + 1 : lease.failureCount;
const status = error instanceof SyncError ? error.status : undefined;
const pause = !signal.aborted && status !== undefined && !retryableStatus(status);
- const delay = failed
- ? retryDelay({ attempt: failureCount, retryMs: timing.retryMs })
- : timing.retryMs;
log({
code,
ownerId: lease.ownerId,
@@ -68,8 +63,7 @@ export class AcquisitionService {
fields: {
...(error instanceof SyncError ? error.diagnostics : {}),
...(status === undefined ? {} : { httpStatus: status }),
- failureCount,
- ...(pause ? { paused: true } : { retryAfterMs: delay }),
+ ...(pause ? { paused: true } : {}),
},
});
this.finish({
@@ -81,8 +75,6 @@ export class AcquisitionService {
? 'interrupted'
: 'retrying',
errorCode: code === 'waiting_for_capacity' ? undefined : code,
- delay,
- failureCount,
pause,
});
}
diff --git a/packages/sync/test/acquisition-retry.test.ts b/packages/sync/test/acquisition-retry.test.ts
index 5646faa..bcb949b 100644
--- a/packages/sync/test/acquisition-retry.test.ts
+++ b/packages/sync/test/acquisition-retry.test.ts
@@ -1,4 +1,4 @@
-import { expect, spyOn, test } from 'bun:test';
+import { expect, test } from 'bun:test';
import type { SyncEvent } from '../src/execution/diagnostics';
import type { SyncContext } from '../src/models/definition';
import { SyncError } from '../src/models/error';
@@ -11,7 +11,6 @@ import { AcquisitionService } from '../src/services/acquisition';
import {
accepted,
alpha,
- beta,
configure,
fixture,
page,
@@ -20,29 +19,14 @@ import {
storage,
} from './support';
-const retryMs = 30_000;
-const maxBackoff = 3_600_000;
-
test.each(['paused', 'interrupted', 'timed_out', 'waiting_for_capacity'])(
- '%s distinguishes source failures from intentional waits',
+ '%s waits for another poll instead of retrying during the drain',
async (state) => {
const f = repositories();
- const previousFailures = 2;
const events: SyncEvent[] = [];
- const prior = f.acquisition.claim(defaultTiming.leaseMs)!;
- f.acquisition.finish({
- lease: prior,
- state: 'retrying',
- errorCode: 'connector_request_failed',
- delay: 0,
- failureCount: previousFailures,
- });
const service = new AcquisitionService({
repository: f.acquisition,
- assets: new SqliteAssets({
- db: f.db,
- maxBytes: defaultLimits.maxPendingAssetBytes,
- }),
+ assets: new SqliteAssets({ db: f.db, maxBytes: defaultLimits.maxPendingAssetBytes }),
files: new DirectoryAssets(`${f.files.path}.assets`),
maxAssetBytes: defaultLimits.maxAssetBytes,
registry: new Registry({
@@ -67,125 +51,29 @@ test.each(['paused', 'interrupted', 'timed_out', 'waiting_for_capacity'])(
state === 'timed_out' ? new DOMException('deadline', 'TimeoutError') : state,
);
await service.execute({ lease: service.claim()!, signal });
- const failureDelay = 120_000;
- expect(events[0]).toMatchObject({
- code: state,
- fields: {
- failureCount: state === 'timed_out' ? previousFailures + 1 : previousFailures,
- retryAfterMs: state === 'timed_out' ? failureDelay : retryMs,
- },
- });
+ expect(events[0]?.code).toBe(state);
+ expect(f.catalog.sync({ ...alpha, id: f.sync.id }).status).toBe(
+ state === 'timed_out'
+ ? 'retrying'
+ : state === 'waiting_for_capacity'
+ ? state
+ : 'interrupted',
+ );
+ expect(service.claim()).toBeUndefined();
+ service.poll();
+ expect(service.claim()?.sync.id).toBe(f.sync.id);
} finally {
f.close();
}
},
);
-test('source failures back off durably despite partial progress, then reset on success', async () => {
- const files = storage();
- const events: SyncEvent[] = [];
- let now = Date.now();
- const clock = spyOn(Date, 'now').mockImplementation(() => now);
- let failing = true;
- const options = {
- databasePath: files.path,
- definitions: [
- {
- ...fixture,
- load: () => ({
- // biome-ignore lint/suspicious/useAwait: The fixture implements the asynchronous execution boundary.
- async step({ checkpoint }: SyncContext) {
- if (checkpoint === 0) {
- return page;
- }
- if (failing) {
- throw new Error('private upstream payload');
- }
- return { ...page, complete: true };
- },
- }),
- },
- ],
- destinationTypes: { local: accepted },
- timing: { retryMs },
- onEvent: (event: SyncEvent) => events.push(event),
- };
- let engine = createSyncRuntime(options);
- try {
- const sync = await configure(engine);
- const scope = { ...alpha, id: sync.id };
- const delays = Object.values({
- first: 30_000,
- second: 60_000,
- third: 120_000,
- fourth: 240_000,
- fifth: 480_000,
- sixth: 960_000,
- seventh: 1_920_000,
- eighth: maxBackoff,
- ninth: maxBackoff,
- });
- await engine.tick();
- for (const delay of delays) {
- await engine.tick();
- expect(savedSync({ path: files.path, scope: scope })).toMatchObject({
- checkpoint: 1,
- status: 'retrying',
- errorCode: 'execution_failed',
- retryAt: now + delay,
- });
- expect(events.at(-1)?.fields?.retryAfterMs).toBe(delay);
- const attempts = events.length;
- await engine.tick();
- expect(events).toHaveLength(attempts);
- now += delay - 1;
- await engine.tick();
- expect(events).toHaveLength(attempts);
- now++;
- // Reopen the actual SQLite database on each attempt, to verify persisted backoff.
- await engine.close();
- engine = createSyncRuntime(options);
- }
- expect(events.at(-1)?.fields?.failureCount).toBe(delays.length);
- expect(JSON.stringify(events)).not.toContain('private upstream payload');
-
- // A different owner starts at the base delay, even while the first source is backed off.
- const destination = { type: 'local', input: {} };
- const other = await engine.api.createSync({
- ...beta,
- definition: fixture.definition.id,
- config: { count: 1 },
- destination,
- });
- await engine.tick();
- await engine.tick();
- await engine.tick();
- expect(engine.api.sync({ ...beta, id: other.id }).retryAt).toBe(now + retryMs);
- await engine.api.setEnabled({ ...beta, id: other.id, enabled: false });
-
- failing = false;
- engine.api.runNow(scope);
- await engine.tick();
- expect(engine.api.sync(scope).status).toBe('succeeded');
- failing = true;
- engine.api.runNow(scope);
- await engine.tick();
- expect(engine.api.sync(scope).retryAt).toBe(now + retryMs);
- expect(events.at(-1)?.fields?.failureCount).toBe(1);
- } finally {
- await engine.close();
- clock.mockRestore();
- files.close();
- }
-});
-
test.each(['records', 'assets'])(
- 'checkpoint yields preserve %s backoff across restarts',
+ '%s failures preserve the checkpoint across restarts and resume on the next cron or manual run',
async (mode) => {
const files = storage();
- let now = Date.now();
- const clock = spyOn(Date, 'now').mockImplementation(() => now);
const events: SyncEvent[] = [];
+ const seen: unknown[] = [];
let failing = true;
const options = {
databasePath: files.path,
@@ -194,6 +82,10 @@ test.each(['records', 'assets'])(
...fixture,
load: () => ({
async step(context: SyncContext) {
+ seen.push(context.checkpoint);
+ if (context.checkpoint === 0) {
+ return page;
+ }
if (failing) {
const failure = new SyncError({
code: 'connector_request_failed',
@@ -212,43 +104,43 @@ test.each(['records', 'assets'])(
throw failure;
}
}
- return page;
+ return { ...page, complete: true };
},
}),
},
],
destinationTypes: { local: accepted },
- timing: { retryMs },
onEvent: (event: SyncEvent) => events.push(event),
};
let engine = createSyncRuntime(options);
try {
const sync = await configure(engine);
const scope = { ...alpha, id: sync.id };
+ await engine.runDue();
+ expect(seen).toEqual([0, 1]);
+ expect(savedSync({ path: files.path, scope })).toMatchObject({
+ checkpoint: 1,
+ status: 'retrying',
+ });
await engine.tick();
- expect(engine.api.sync(scope).retryAt).toBe(now + retryMs);
- expect(events.at(-1)?.fields).toMatchObject({ httpStatus: 429, retryAfterMs: retryMs });
- expect(JSON.stringify(events)).not.toContain('private upstream payload');
- now += retryMs;
+ expect(seen).toEqual([0, 1]);
+ await engine.close();
+ engine = createSyncRuntime(options);
await engine.tick();
- const secondDelay = 60_000;
- expect(engine.api.sync(scope).retryAt).toBe(now + secondDelay);
- now += secondDelay;
+ expect(seen).toEqual([0, 1]);
+ await engine.runDue();
+ expect(seen).toEqual([0, 1, 1]);
+ expect(events.at(-1)?.fields).toEqual({ httpStatus: 429 });
+ expect(JSON.stringify(events)).not.toContain('private upstream payload');
failing = false;
+ engine.api.runNow(scope);
await engine.tick();
- expect(savedSync({ path: files.path, scope: scope })).toMatchObject({
- status: 'ready',
- retryAt: null,
- });
- await engine.close();
- engine = createSyncRuntime(options);
- failing = true;
+ expect(seen).toEqual([0, 1, 1, 1]);
+ expect(engine.api.sync(scope).status).toBe('succeeded');
await engine.tick();
- const thirdDelay = 120_000;
- expect(engine.api.sync(scope).retryAt).toBe(now + thirdDelay);
+ expect(seen).toEqual([0, 1, 1, 1]);
} finally {
await engine.close();
- clock.mockRestore();
files.close();
}
},
diff --git a/packages/sync/test/capacity.test.ts b/packages/sync/test/capacity.test.ts
index fb579c5..575a163 100644
--- a/packages/sync/test/capacity.test.ts
+++ b/packages/sync/test/capacity.test.ts
@@ -70,7 +70,7 @@ test('download reservations account for other in-flight downloads and survive re
'waiting for capacity',
);
// Only deletion of the actual file releases its charge, not completion of the source run.
- f.acquisition.finish({ lease: first, state: 'interrupted', delay: leaseMs });
+ f.acquisition.finish({ lease: first, state: 'interrupted' });
expect(() => restored.reserve({ lease: second, id: right, bytes: 5 })).toThrow(
'waiting for capacity',
);
diff --git a/packages/sync/test/connector-failure.test.ts b/packages/sync/test/connector-failure.test.ts
index fefe882..a548b83 100644
--- a/packages/sync/test/connector-failure.test.ts
+++ b/packages/sync/test/connector-failure.test.ts
@@ -1,4 +1,4 @@
-import { expect, spyOn, test } from 'bun:test';
+import { expect, test } from 'bun:test';
import { createConnectorClient } from '../src/connector/client';
import { connectorFailure } from '../src/connector/failure';
import type { SyncEvent } from '../src/execution/diagnostics';
@@ -77,11 +77,9 @@ function client(response: () => Promise) {
});
}
-test('connector failures use engine backoff despite timing headers and retain only safe scoped diagnostics', async () => {
+test('connector failures retry on the next poll despite timing headers and retain safe scoped diagnostics', async () => {
const files = storage();
const events: SyncEvent[] = [];
- let now = Date.now();
- const clock = spyOn(Date, 'now').mockImplementation(() => now);
let retryAfter = '120';
const connector = client(() =>
Promise.resolve(
@@ -137,36 +135,28 @@ test('connector failures use engine backoff despite timing headers and retain on
connectorErrorCode: 'rate_limited',
providerStatus: rateLimited,
httpStatus: rateLimited,
- failureCount: 1,
- retryAfterMs: 30_000,
},
},
]);
expect(JSON.stringify(events)).not.toContain(secret);
const scope = { ...alpha, id: sync.id };
- const firstDelay = 30_000;
- const secondDelay = 60_000;
expect(savedSync({ path: files.path, scope: scope })).toMatchObject({
checkpoint: 0,
- retryAt: now + firstDelay,
+ status: 'retrying',
});
- now += firstDelay;
const providerDelay = 120_000;
- retryAfter = new Date(now + providerDelay).toUTCString();
- await engine.tick();
+ retryAfter = new Date(Date.now() + providerDelay).toUTCString();
+ await engine.runDue();
expect(savedSync({ path: files.path, scope: scope })).toMatchObject({
checkpoint: 0,
- retryAt: now + secondDelay,
+ status: 'retrying',
});
expect(events.at(-1)?.fields).toMatchObject({
httpStatus: rateLimited,
- failureCount: 2,
- retryAfterMs: secondDelay,
});
expect(JSON.stringify(events)).not.toContain(secret);
} finally {
await engine.close();
- clock.mockRestore();
files.close();
}
});
diff --git a/packages/sync/test/shared-backoff.test.ts b/packages/sync/test/delivery-backoff.test.ts
similarity index 69%
rename from packages/sync/test/shared-backoff.test.ts
rename to packages/sync/test/delivery-backoff.test.ts
index 9d01dee..cb8ac12 100644
--- a/packages/sync/test/shared-backoff.test.ts
+++ b/packages/sync/test/delivery-backoff.test.ts
@@ -1,23 +1,19 @@
import { expect, spyOn, test } from 'bun:test';
-import type { SyncContext } from '../src/models/definition';
import { createSyncRuntime } from '../src/runtime';
import { accepted, alpha, configure, fixture, page, storage } from './support';
-test('source and delivery failures use the same durable exponential backoff', async () => {
+test('delivery failures retain exponential backoff independently of source polling', async () => {
const files = storage();
let now = Date.now();
const clock = spyOn(Date, 'now').mockImplementation(() => now);
- let failingSource = '';
const engine = createSyncRuntime({
databasePath: files.path,
definitions: [
{
...fixture,
load: () => ({
- step(context: SyncContext) {
- return context.syncId === failingSource
- ? Promise.reject(new Error('temporary'))
- : Promise.resolve({ ...page, complete: true });
+ step() {
+ return Promise.resolve({ ...page, complete: true });
},
}),
},
@@ -27,8 +23,6 @@ test('source and delivery failures use the same durable exponential backoff', as
},
});
try {
- const source = await configure(engine);
- failingSource = source.id;
const destinationSync = await configure(engine);
await engine.tick();
await engine.tick();
@@ -45,13 +39,13 @@ test('source and delivery failures use the same durable exponential backoff', as
});
for (const delay of delays) {
const due = now + delay;
- expect(engine.api.sync({ ...alpha, id: source.id }).retryAt).toBe(due);
+ await engine.runDue();
expect(
engine.api.deliveries({ ...alpha, syncId: destinationSync.id }).deliveries[0]
?.nextAttemptAt,
).toBe(due);
now = due;
- await engine.tick();
+ await engine.runDue();
}
} finally {
await engine.close();
diff --git a/packages/sync/test/persistence.test.ts b/packages/sync/test/persistence.test.ts
index efb9ef8..3d00cf9 100644
--- a/packages/sync/test/persistence.test.ts
+++ b/packages/sync/test/persistence.test.ts
@@ -13,7 +13,7 @@ test.each([false, true])(
const before = f.catalog.sync({ ...alpha, id: f.sync.id });
const pollsBefore = f.catalog.polls({ ...alpha, id: f.sync.id });
f.db.exec(
- "CREATE TRIGGER fail_checkpoint BEFORE UPDATE OF retry_at ON syncs BEGIN SELECT RAISE(ABORT,'injected'); END;",
+ "CREATE TRIGGER fail_checkpoint BEFORE UPDATE OF status ON syncs BEGIN SELECT RAISE(ABORT,'injected'); END;",
);
expect(() =>
f.acquisition.commit({
diff --git a/packages/sync/test/polling-migration.test.ts b/packages/sync/test/polling-migration.test.ts
index 01326f8..b84c7f7 100644
--- a/packages/sync/test/polling-migration.test.ts
+++ b/packages/sync/test/polling-migration.test.ts
@@ -1,7 +1,7 @@
import { Database } from 'bun:sqlite';
import { expect, test } from 'bun:test';
import { openDatabase } from '../src/db/client';
-import releasedSchema from '../src/db/schema.sql' with { type: 'text' };
+import releasedSchema from '../src/db/migrations/0000-initial.sql' with { type: 'text' };
import type { Deliverable } from '../src/models/delivery';
import { defaultLimits } from '../src/models/limits';
import { SqliteAcquisition } from '../src/repositories/acquisition/sqlite';
@@ -9,7 +9,7 @@ import { SqliteCatalog } from '../src/repositories/catalog/sqlite';
import { SqliteDeliveries } from '../src/repositories/delivery/sqlite';
import { alpha, fixture, storage } from './support';
-const migrationNames = ['0000-initial-schema', '0001-centralize-polling'];
+const migrationNames = ['0000-initial.sql', '0001-centralize-polling.sql'];
const leaseMs = 60_000;
const futureRetry = 9_000_000_000_000;
const retryStates = ['retrying', 'interrupted', 'waiting_for_capacity'];
@@ -93,7 +93,7 @@ function history(db: Database) {
.all();
}
-test('released databases adopt migration history without losing progress, leases, retries or queued data', () => {
+test('released databases adopt migration history without losing progress, leases or queued data', () => {
const files = storage();
let db = legacyDatabase(files.path);
const before = db.query('SELECT * FROM syncs ORDER BY id').all() as Record[];
@@ -103,10 +103,9 @@ test('released databases adopt migration history without losing progress, leases
db = openDatabase(files.path);
const after = db.query('SELECT * FROM syncs ORDER BY id').all();
expect(after).toEqual(
- before.map(({ interval_ms: _interval, next_due_at, ...row }) => ({
- ...row,
- retry_at: retryStates.includes(String(row.status)) ? next_due_at : null,
- })),
+ before.map(
+ ({ interval_ms: _interval, next_due_at: _due, failure_count: _failures, ...row }) => row,
+ ),
);
expect(db.query('PRAGMA foreign_key_check').all()).toEqual([]);
expect(queuedState(db)).toEqual(queue);
@@ -123,8 +122,11 @@ test('released databases adopt migration history without losing progress, leases
expect(acquisition.claim(leaseMs)?.sync.id).toBe('ready');
expect(acquisition.claim(leaseMs)).toBeUndefined();
acquisition.poll();
- expect(acquisition.claim(leaseMs)?.sync.id).toBe('succeeded');
- expect(acquisition.claim(leaseMs)).toBeUndefined();
+ const claimed = [];
+ for (let lease = acquisition.claim(leaseMs); lease; lease = acquisition.claim(leaseMs)) {
+ claimed.push(lease.sync.id);
+ }
+ expect(claimed.sort()).toEqual(['succeeded', ...retryStates].sort());
expect(new SqliteDeliveries(db).claim(leaseMs)?.delivery).toEqual(delivery);
expect(
catalog.createSync({
@@ -133,8 +135,8 @@ test('released databases adopt migration history without losing progress, leases
config: { count: 1 },
destination: { type: 'local', config: {} },
initialCheckpoint: 0,
- }).retryAt,
- ).toBeNull();
+ }).status,
+ ).toBe('ready');
} finally {
db.close();
files.close();
@@ -148,7 +150,8 @@ test('fresh databases apply the same ordered migration history', () => {
.query<{ name: string }, []>("SELECT name FROM pragma_table_info('syncs')")
.all()
.map(({ name }) => name);
- expect(columns).toContain('retry_at');
+ expect(columns).not.toContain('retry_at');
+ expect(columns).not.toContain('failure_count');
expect(columns).not.toContain('interval_ms');
expect(columns).not.toContain('next_due_at');
});
@@ -165,16 +168,14 @@ test.each([false, true])(
'2026-10-01T00:00:00.000Z',
);
}
- db.exec(
- "CREATE TRIGGER fail_migration BEFORE UPDATE ON syncs BEGIN SELECT RAISE(ABORT,'injected'); END",
- );
+ db.exec('CREATE INDEX fail_migration ON syncs(next_due_at)');
const before = db.query('SELECT * FROM syncs ORDER BY id').all();
const schemaBefore = db.query('SELECT * FROM sqlite_schema ORDER BY name').all();
const queue = queuedState(db);
const historyBefore = recorded ? history(db) : [];
db.close();
try {
- expect(() => openDatabase(files.path)).toThrow('injected');
+ expect(() => openDatabase(files.path)).toThrow('fail_migration');
db = new Database(files.path);
expect(db.query('SELECT * FROM syncs ORDER BY id').all()).toEqual(before);
expect(db.query('SELECT * FROM sqlite_schema ORDER BY name').all()).toEqual(schemaBefore);
@@ -182,7 +183,7 @@ test.each([false, true])(
if (recorded) {
expect(history(db)).toEqual(historyBefore);
}
- db.exec('DROP TRIGGER fail_migration');
+ db.exec('DROP INDEX fail_migration');
db.close();
db = openDatabase(files.path);
expect(history(db).map(({ name }) => name)).toEqual(migrationNames);
diff --git a/packages/sync/test/run-due.test.ts b/packages/sync/test/run-due.test.ts
index 66091b0..5220603 100644
--- a/packages/sync/test/run-due.test.ts
+++ b/packages/sync/test/run-due.test.ts
@@ -1,4 +1,4 @@
-import { expect, spyOn, test } from 'bun:test';
+import { expect, test } from 'bun:test';
import { createSyncRuntime } from '../src/runtime';
import { isolatedScheduler, runRegisteredCron } from './cron-support';
import { accepted, alpha, beta, configure, fixture, runtime, storage } from './support';
@@ -199,9 +199,7 @@ async function until(check: () => boolean) {
}
}
-test('each cron round polls healthy syncs without bypassing source retry backoff', async () => {
- let now = Date.now();
- const clock = spyOn(Date, 'now').mockImplementation(() => now);
+test('each cron round attempts healthy and failed syncs exactly once', async () => {
const calls: number[] = [];
let failing = true;
const f = runtime({
@@ -232,16 +230,16 @@ test('each cron round polls healthy syncs without bypassing source retry backoff
await f.engine.runDue();
await f.engine.runDue();
expect(calls.filter((count) => count === 1)).toHaveLength(2);
- expect(calls.filter((count) => count === 2)).toHaveLength(1);
+ expect(calls.filter((count) => count === 2)).toHaveLength(2);
const scope = { ...alpha, id: syncs[1]!.id };
- expect(f.engine.api.sync(scope).retryAt).toBe(now + f.options.timing.retryMs);
- now += f.options.timing.retryMs;
+ expect(f.engine.api.sync(scope).status).toBe('retrying');
+ await f.engine.tick();
+ expect(calls.filter((count) => count === 2)).toHaveLength(2);
failing = false;
await f.engine.runDue();
- expect(calls.filter((count) => count === 2)).toHaveLength(2);
- expect(f.engine.api.sync(scope)).toMatchObject({ status: 'succeeded', retryAt: null });
+ expect(calls.filter((count) => count === 2)).toEqual([2, 2, 2]);
+ expect(f.engine.api.sync(scope).status).toBe('succeeded');
} finally {
await f.close();
- clock.mockRestore();
}
});
diff --git a/packages/sync/test/source-http-error.test.ts b/packages/sync/test/source-http-error.test.ts
index 592a591..948b784 100644
--- a/packages/sync/test/source-http-error.test.ts
+++ b/packages/sync/test/source-http-error.test.ts
@@ -1,10 +1,9 @@
-import { expect, spyOn, test } from 'bun:test';
+import { expect, test } from 'bun:test';
import type { SyncEvent } from '../src/execution/diagnostics';
import { SourceHttpError, type SyncContext } from '../src/models/definition';
import { createSyncRuntime } from '../src/runtime';
import { accepted, alpha, configure, fixture, page, savedSync, storage } from './support';
-const retryMs = 30_000;
const rateLimited = 429;
const forbidden = 403;
const unavailable = 503;
@@ -20,11 +19,9 @@ test.each(
unavailable: 503,
gatewayTimeout: 504,
}),
-)('HTTP %s backs off at the shared engine boundary', async (status) => {
+)('HTTP %s retries on the next polling round', async (status) => {
const files = storage();
const events: SyncEvent[] = [];
- const now = Date.now();
- const clock = spyOn(Date, 'now').mockReturnValue(now);
const engine = createSyncRuntime({
databasePath: files.path,
definitions: [
@@ -46,16 +43,16 @@ test.each(
checkpoint: 0,
status: 'retrying',
errorCode: `source_http_${status}`,
- retryAt: now + retryMs,
});
expect(events[0]?.fields).toEqual({
httpStatus: status,
- failureCount: 1,
- retryAfterMs: retryMs,
});
+ await engine.tick();
+ expect(events).toHaveLength(1);
+ await engine.runDue();
+ expect(events).toHaveLength(2);
} finally {
await engine.close();
- clock.mockRestore();
files.close();
}
});
@@ -108,7 +105,7 @@ test.each(
});
await engine.close();
engine = createSyncRuntime(options);
- await engine.tick();
+ await engine.runDue();
expect(seen).toEqual([0, 1]);
rejected = false;
await engine.api.setEnabled({ ...scope, enabled: true });
@@ -129,8 +126,6 @@ test.each([rateLimited, forbidden, unavailable])(
async (status) => {
const files = storage();
const events: SyncEvent[] = [];
- let now = Date.now();
- const clock = spyOn(Date, 'now').mockImplementation(() => now);
let rejected = true;
const options = {
databasePath: files.path,
@@ -165,38 +160,34 @@ test.each([rateLimited, forbidden, unavailable])(
try {
const sync = await configure(engine);
const scope = { ...alpha, id: sync.id };
- const delays = Object.values({ first: 30_000, second: 60_000, third: 120_000 });
+ const rounds = 3;
await engine.tick();
- for (const [index, delay] of delays.entries()) {
- await engine.tick();
+ for (let round = 0; round < rounds; round++) {
+ await engine.runDue();
expect(savedSync({ path: files.path, scope: scope })).toMatchObject({
enabled: status !== forbidden,
checkpoint: 1,
status: status === forbidden ? 'disabled' : 'retrying',
errorCode: `source_http_${status}`,
- retryAt: now + delay,
});
expect(events.at(-1)?.fields).toEqual({
httpStatus: status,
- failureCount: index + 1,
- ...(status === forbidden ? { paused: true } : { retryAfterMs: delay }),
+ ...(status === forbidden ? { paused: true } : {}),
});
await engine.close();
engine = createSyncRuntime(options);
- now += delay;
if (status === forbidden) {
await engine.api.setEnabled({ ...scope, enabled: true });
}
}
rejected = false;
- await engine.tick();
+ await engine.runDue();
expect(savedSync({ path: files.path, scope: scope })).toMatchObject({
enabled: true,
status: 'succeeded',
});
} finally {
await engine.close();
- clock.mockRestore();
files.close();
}
},
From e144c2cf55df2e54f1d91980a478a07e37acc7e2 Mon Sep 17 00:00:00 2001
From: massimoalbarello
Date: Fri, 2 Oct 2026 21:50:02 +0100
Subject: [PATCH 4/5] fix: migrate existing hosts with plain SQL
---
bun.lock | 3 +
packages/sync/package.json | 1 +
packages/sync/src/db/client.ts | 13 +-
packages/sync/src/db/migrate.ts | 31 +--
.../sync/src/db/migrations/0000-initial.sql | 10 -
.../db/migrations/0001-centralize-polling.sql | 2 -
packages/sync/test/polling-migration.test.ts | 14 ++
packages/sync/test/released-upgrade.test.ts | 192 ++++++++++++++++++
8 files changed, 236 insertions(+), 30 deletions(-)
create mode 100644 packages/sync/test/released-upgrade.test.ts
diff --git a/bun.lock b/bun.lock
index 454f0db..a632d6e 100644
--- a/bun.lock
+++ b/bun.lock
@@ -107,6 +107,7 @@
"proper-lockfile": "4.1.2",
},
"devDependencies": {
+ "open-sync-previous": "https://registry.npmjs.org/@context-use/open-sync/-/open-sync-0.3.1.tgz",
"typescript": "catalog:",
},
},
@@ -1255,6 +1256,8 @@
"once": ["once@1.4.0", "", { "dependencies": { "wrappy": "1" } }, "sha512-lNaJgI+2Q5URQBkccEKHTQOPaXdUxnZZElQTZY0MFUAuaEqe1E+Nyvgdz/aIyNi6Z9MzO5dv1H8n58/GELp3+w=="],
+ "open-sync-previous": ["@context-use/open-sync@https://registry.npmjs.org/@context-use/open-sync/-/open-sync-0.3.1.tgz", { "dependencies": { "@cfworker/json-schema": "4.1.1", "@oomol-lab/open-connector": "1.7.0", "elysia": "^1.4.30", "minisearch": "^7.2.0" } }, "sha512-RONqbwuPxb6ekvIssRlRrl3CHj1Npn+bcD6Q4EcCK3pialOyvO6LNg2nMf79TkCNEjtu02bOeaKHU1CMWqd2wQ=="],
+
"openapi-types": ["openapi-types@12.1.3", "", {}, "sha512-N4YtSYJqghVu4iek2ZUvcN/0aqH1kRDuNqzcycDxhOUpg7GdvLa2F3DgS6yBNhInhv2r/6I0Flkn7CqL8+nIcw=="],
"optionator": ["optionator@0.9.4", "", { "dependencies": { "deep-is": "^0.1.3", "fast-levenshtein": "^2.0.6", "levn": "^0.4.1", "prelude-ls": "^1.2.1", "type-check": "^0.4.0", "word-wrap": "^1.2.5" } }, "sha512-6IpQ7mKUxRcZNLIObR0hz7lxsapSSIYNZJwXPGeF0mTVqGKFIXj1DQcMoT22S3ROcLyY/rz0PWaWZ9ayWmad9g=="],
diff --git a/packages/sync/package.json b/packages/sync/package.json
index 752ed5c..1aacf53 100644
--- a/packages/sync/package.json
+++ b/packages/sync/package.json
@@ -32,6 +32,7 @@
"proper-lockfile": "4.1.2"
},
"devDependencies": {
+ "open-sync-previous": "https://registry.npmjs.org/@context-use/open-sync/-/open-sync-0.3.1.tgz",
"typescript": "catalog:"
},
"engines": {
diff --git a/packages/sync/src/db/client.ts b/packages/sync/src/db/client.ts
index 9813eb6..8cfbd8d 100644
--- a/packages/sync/src/db/client.ts
+++ b/packages/sync/src/db/client.ts
@@ -1,11 +1,18 @@
import { Database } from 'bun:sqlite';
+import { DatabaseSync } from 'node:sqlite';
import { runMigrations } from './migrate';
export function openDatabase(path: string): Database {
- const db = new Database(path, { create: true, strict: true });
+ // Keep in-memory databases private while both startup connections share the same database.
+ const filename =
+ path === ':memory:' ? `file:open-sync-${crypto.randomUUID()}?mode=memory&cache=shared` : path;
+ // Unlike bun:sqlite in Bun 1.4.0, this driver propagates failures inside SQL batches.
+ using migration = new DatabaseSync(filename);
+ migration.exec('PRAGMA busy_timeout=5000; PRAGMA journal_mode=WAL;');
+ runMigrations(migration);
+ const db = new Database(filename, { create: true, strict: true });
try {
- db.exec('PRAGMA foreign_keys=ON; PRAGMA journal_mode=WAL; PRAGMA busy_timeout=5000;');
- runMigrations(db);
+ db.exec('PRAGMA foreign_keys=ON; PRAGMA busy_timeout=5000;');
return db;
} catch (error) {
db.close();
diff --git a/packages/sync/src/db/migrate.ts b/packages/sync/src/db/migrate.ts
index cf15efc..f2d0680 100644
--- a/packages/sync/src/db/migrate.ts
+++ b/packages/sync/src/db/migrate.ts
@@ -1,4 +1,4 @@
-import type { Database } from 'bun:sqlite';
+import type { DatabaseSync } from 'node:sqlite';
import initialSchema from './migrations/0000-initial.sql' with { type: 'text' };
import centralizePolling from './migrations/0001-centralize-polling.sql' with { type: 'text' };
@@ -8,29 +8,30 @@ const migrations = [
{ name: '0001-centralize-polling.sql', sql: centralizePolling },
] as const;
-export function runMigrations(db: Database): void {
- db.transaction(() => {
- db.query(`CREATE TABLE IF NOT EXISTS __migrations (
+export function runMigrations(db: DatabaseSync): void {
+ db.exec('BEGIN IMMEDIATE');
+ try {
+ db.exec(`CREATE TABLE IF NOT EXISTS __migrations (
name TEXT PRIMARY KEY, applied_at TEXT NOT NULL
- )`).run();
- const applied = db
- .query<{ name: string }, []>('SELECT name FROM __migrations ORDER BY name')
- .all();
+ )`);
+ const applied = db.prepare('SELECT name FROM __migrations ORDER BY name').all();
for (const [index, { name }] of applied.entries()) {
if (name !== migrations[index]?.name) {
throw new Error('Unsupported sync database migration history');
}
}
for (const migration of migrations.slice(applied.length)) {
- // Explicit boundaries also support triggers and strings containing semicolons.
- // Executing each statement separately surfaces errors hidden by Bun 1.4 SQL batches.
- for (const statement of migration.sql.split('--> statement-breakpoint')) {
- db.query(statement).run();
- }
- db.query('INSERT INTO __migrations (name,applied_at) VALUES (?,?)').run(
+ db.exec(migration.sql);
+ db.prepare('INSERT INTO __migrations (name,applied_at) VALUES (?,?)').run(
migration.name,
new Date().toISOString(),
);
}
- }).immediate();
+ db.exec('COMMIT');
+ } catch (error) {
+ if (db.isTransaction) {
+ db.exec('ROLLBACK');
+ }
+ throw error;
+ }
}
diff --git a/packages/sync/src/db/migrations/0000-initial.sql b/packages/sync/src/db/migrations/0000-initial.sql
index 2dbd1b3..7a288f1 100644
--- a/packages/sync/src/db/migrations/0000-initial.sql
+++ b/packages/sync/src/db/migrations/0000-initial.sql
@@ -9,14 +9,12 @@ CREATE TABLE IF NOT EXISTS syncs (
failure_count INTEGER NOT NULL DEFAULT 0, resync INTEGER NOT NULL DEFAULT 0,
PRIMARY KEY(owner_id, id)
);
---> statement-breakpoint
CREATE TABLE IF NOT EXISTS record_state (
owner_id TEXT NOT NULL, sync_id TEXT NOT NULL, kind TEXT NOT NULL, id TEXT NOT NULL,
hash TEXT NOT NULL, revision INTEGER NOT NULL, deleted INTEGER NOT NULL,
PRIMARY KEY(owner_id, sync_id, kind, id),
FOREIGN KEY(owner_id, sync_id) REFERENCES syncs(owner_id, id)
);
---> statement-breakpoint
-- Polling diagnostics only; execution leases and scheduling belong to syncs.
CREATE TABLE IF NOT EXISTS sync_polls (
id INTEGER PRIMARY KEY AUTOINCREMENT, owner_id TEXT NOT NULL, sync_id TEXT NOT NULL,
@@ -25,11 +23,8 @@ CREATE TABLE IF NOT EXISTS sync_polls (
state TEXT NOT NULL, error_code TEXT,
FOREIGN KEY(owner_id, sync_id) REFERENCES syncs(owner_id, id) ON DELETE CASCADE
);
---> statement-breakpoint
CREATE INDEX IF NOT EXISTS sync_poll_history ON sync_polls(owner_id,sync_id,id);
---> statement-breakpoint
CREATE UNIQUE INDEX IF NOT EXISTS sync_poll_active ON sync_polls(owner_id,sync_id) WHERE completed_at IS NULL;
---> statement-breakpoint
CREATE TABLE IF NOT EXISTS deliveries (
sequence INTEGER PRIMARY KEY AUTOINCREMENT, owner_id TEXT NOT NULL, id TEXT NOT NULL,
sync_id TEXT NOT NULL, body TEXT NOT NULL,
@@ -39,11 +34,8 @@ CREATE TABLE IF NOT EXISTS deliveries (
UNIQUE(owner_id, id),
FOREIGN KEY(owner_id, sync_id) REFERENCES syncs(owner_id, id)
);
---> statement-breakpoint
CREATE INDEX IF NOT EXISTS delivery_order ON deliveries(owner_id,sync_id,sequence);
---> statement-breakpoint
CREATE INDEX IF NOT EXISTS deliveries_due ON deliveries(state, due_at);
---> statement-breakpoint
-- Staging, queued bytes, and pending cleanup share one delivery-owned ledger.
-- No foreign keys: cleanup must survive deletion of the delivery or sync.
CREATE TABLE IF NOT EXISTS delivery_assets (
@@ -54,7 +46,5 @@ CREATE TABLE IF NOT EXISTS delivery_assets (
bytes INTEGER NOT NULL DEFAULT 0 CHECK(bytes>=0), delivery_id TEXT,
UNIQUE(owner_id,sync_id,generation,asset_id,asset_version)
);
---> statement-breakpoint
CREATE INDEX IF NOT EXISTS delivery_asset_queue ON delivery_assets(owner_id,delivery_id);
---> statement-breakpoint
CREATE INDEX IF NOT EXISTS delivery_asset_sync ON delivery_assets(owner_id,sync_id);
diff --git a/packages/sync/src/db/migrations/0001-centralize-polling.sql b/packages/sync/src/db/migrations/0001-centralize-polling.sql
index 2eed189..a88ce4d 100644
--- a/packages/sync/src/db/migrations/0001-centralize-polling.sql
+++ b/packages/sync/src/db/migrations/0001-centralize-polling.sql
@@ -1,5 +1,3 @@
ALTER TABLE syncs DROP COLUMN interval_ms;
---> statement-breakpoint
ALTER TABLE syncs DROP COLUMN next_due_at;
---> statement-breakpoint
ALTER TABLE syncs DROP COLUMN failure_count;
diff --git a/packages/sync/test/polling-migration.test.ts b/packages/sync/test/polling-migration.test.ts
index b84c7f7..65e6acc 100644
--- a/packages/sync/test/polling-migration.test.ts
+++ b/packages/sync/test/polling-migration.test.ts
@@ -156,6 +156,20 @@ test('fresh databases apply the same ordered migration history', () => {
expect(columns).not.toContain('next_due_at');
});
+test('in-memory databases retain migrated tables and stay isolated after startup', () => {
+ using first = openDatabase(':memory:');
+ using second = openDatabase(':memory:');
+ const sync = new SqliteCatalog(first).createSync({
+ ...alpha,
+ definition: fixture.definition.id,
+ config: { count: 1 },
+ destination: { type: 'local', config: {} },
+ initialCheckpoint: 0,
+ });
+ expect(new SqliteCatalog(first).sync({ ...alpha, id: sync.id })).toEqual(sync);
+ expect(new SqliteCatalog(second).syncs(alpha)).toEqual([]);
+});
+
test.each([false, true])(
'failed migration rolls back schema, data and history (baseline recorded: %s)',
(recorded) => {
diff --git a/packages/sync/test/released-upgrade.test.ts b/packages/sync/test/released-upgrade.test.ts
new file mode 100644
index 0000000..d9bf7d1
--- /dev/null
+++ b/packages/sync/test/released-upgrade.test.ts
@@ -0,0 +1,192 @@
+import { Database } from 'bun:sqlite';
+import { expect, test } from 'bun:test';
+import { readFile } from 'node:fs/promises';
+import { join } from 'node:path';
+import { createOpenSync as createPreviousOpenSync } from 'open-sync-previous';
+import type { SyncRegistration } from '../src/models/definition';
+import type { Deliverable, DeliveryResult } from '../src/models/delivery';
+import { createOpenSync, type OpenSyncRuntime } from '../src/open-sync';
+import { accepted, alpha, beta, fixture, storage } from './support';
+
+const assetBody = 'attachment created by the released host';
+const longInterval = 3_600_000;
+
+// Create the data with the published package, not a reconstructed copy of its SQL schema.
+// The registry tarball is pinned because an npm alias would resolve to this workspace package.
+test('a host using published 0.3.1 upgrades its existing directory without resetting state', async () => {
+ const files = storage();
+ const dataDirectory = join(files.dir, 'open-sync');
+ const databasePath = join(dataDirectory, 'sync.db');
+ const resumed: unknown[] = [];
+ const received: Array> = [];
+ const registration: SyncRegistration = {
+ ...fixture,
+ load: () => ({
+ async step(context) {
+ if (context.checkpoint !== 0) {
+ throw new Error('temporary source failure');
+ }
+ const file = await context.assets.capture({
+ id: 'attachment',
+ version: '1',
+ name: 'saved.txt',
+ mediaType: 'text/plain',
+ read: () => Promise.resolve(new Blob([assetBody]).stream()),
+ });
+ return {
+ checkpoint: 1,
+ complete: false,
+ records: [
+ {
+ operation: 'upsert',
+ kind: 'item',
+ id: 'saved',
+ data: { value: 1 },
+ assetRefs: { file },
+ },
+ ],
+ };
+ },
+ }),
+ };
+ const host = {
+ dataDirectory,
+ publicUrl: 'http://host/api/open-sync',
+ authorize: () => alpha,
+ canConfigureProviders: () => Promise.resolve(true),
+ };
+ let previous: Awaited> | undefined;
+ let current: OpenSyncRuntime | undefined;
+ try {
+ // The embedding application's own database is outside Open Sync's ownership.
+ using hostDb = new Database(join(files.dir, 'app.db'));
+ hostDb.exec('CREATE TABLE host_records(id TEXT PRIMARY KEY)');
+ hostDb.query('INSERT INTO host_records VALUES (?)').run('existing-host-data');
+ previous = await createPreviousOpenSync({
+ ...host,
+ definitions: [registration],
+ destinationTypes: {
+ local: {
+ ...accepted,
+ deliver: (): Promise =>
+ Promise.resolve({ status: 'rejected', code: 'receiver_unavailable' }),
+ },
+ },
+ });
+ await previous.providers.configure({
+ ...alpha,
+ service: 'github',
+ values: { clientId: 'saved-client', clientSecret: 'saved-secret' },
+ });
+ const sync = await previous.api.createSync({
+ ...alpha,
+ definition: fixture.definition.id,
+ config: { count: 1 },
+ destination: { type: 'local', input: {} },
+ intervalMs: longInterval,
+ });
+ const paused = await previous.api.createSync({
+ ...beta,
+ definition: fixture.definition.id,
+ config: { count: 1 },
+ destination: { type: 'local', input: {} },
+ enabled: false,
+ intervalMs: longInterval,
+ });
+ const scope = { ...alpha, id: sync.id };
+ previous.start();
+ await until(
+ () =>
+ previous!.api.sync(scope).status === 'retrying' &&
+ previous!.api.deliveries({ ...alpha, syncId: sync.id }).deliveries[0]?.state === 'blocked',
+ );
+ const queued = previous.api.deliveries({ ...alpha, syncId: sync.id }).deliveries[0]!;
+ const deliverable = previous.api.deliverable({ ...alpha, syncId: sync.id, id: queued.id });
+ const before = previous.api.sync(scope);
+ const pausedBefore = previous.api.sync({ ...beta, id: paused.id });
+ const polls = previous.api.polls(scope);
+ const providers = await previous.providers.status({ ...alpha, service: 'github' });
+ await previous.close();
+ previous = undefined;
+ const key = await readFile(join(dataDirectory, '.connector-key'));
+ using beforeDb = new Database(databasePath, { readonly: true });
+ expect(
+ beforeDb.query("SELECT name FROM sqlite_schema WHERE name='__migrations'").all(),
+ ).toEqual([]);
+ const records = beforeDb.query('SELECT * FROM record_state').all();
+ const assets = beforeDb.query('SELECT * FROM delivery_assets').all();
+ const upgradedOptions = {
+ ...host,
+ definitions: [
+ {
+ ...registration,
+ load: () => ({
+ step({ checkpoint }: { checkpoint: unknown }) {
+ resumed.push(checkpoint);
+ return Promise.resolve({ checkpoint: 2, records: [], complete: true });
+ },
+ }),
+ },
+ ],
+ destinationTypes: {
+ local: {
+ ...accepted,
+ async deliver({ deliverable }: { deliverable: Deliverable }): Promise {
+ const { openAsset, ...payload } = deliverable;
+ expect(await new Response(await openAsset(deliverable.assets[0]!)).text()).toBe(
+ assetBody,
+ );
+ received.push(payload);
+ return { status: 'accepted' };
+ },
+ },
+ },
+ };
+ current = await createOpenSync(upgradedOptions);
+ const { intervalMs: _interval, nextDueAt: _due, ...retained } = before;
+ const { intervalMs: _pausedInterval, nextDueAt: _pausedDue, ...retainedPaused } = pausedBefore;
+ expect(current.api.sync(scope)).toEqual(retained);
+ expect(current.api.sync({ ...beta, id: paused.id })).toEqual(retainedPaused);
+ expect(() => current!.api.sync({ ...beta, id: sync.id })).toThrow('not found');
+ expect(current.api.polls(scope)).toEqual(polls);
+ expect(current.api.deliverable({ ...alpha, syncId: sync.id, id: queued.id })).toEqual(
+ deliverable,
+ );
+ expect(await current.providers.status({ ...alpha, service: 'github' })).toEqual(providers);
+ expect(await readFile(join(dataDirectory, '.connector-key'))).toEqual(key);
+ using afterDb = new Database(databasePath, { readonly: true });
+ expect(afterDb.query('SELECT * FROM record_state').all()).toEqual(records);
+ expect(afterDb.query('SELECT * FROM delivery_assets').all()).toEqual(assets);
+ expect(afterDb.query('PRAGMA foreign_key_check').all()).toEqual([]);
+ const history = afterDb.query('SELECT * FROM __migrations ORDER BY name').all();
+ expect(history).toHaveLength(2);
+ await current.close();
+ current = await createOpenSync(upgradedOptions);
+ expect(afterDb.query('SELECT * FROM __migrations ORDER BY name').all()).toEqual(history);
+ current.api.retryDelivery({ ...alpha, syncId: sync.id, id: queued.id });
+ await current.runDue();
+ expect(resumed).toEqual([1]);
+ expect(received).toEqual([deliverable]);
+ expect(current.api.sync(scope).status).toBe('succeeded');
+ expect(current.api.sync({ ...beta, id: paused.id }).status).toBe('disabled');
+ expect(current.api.status(alpha).queue.pendingRecords).toBe(0);
+ expect(hostDb.query('SELECT * FROM host_records').all()).toEqual([
+ { id: 'existing-host-data' },
+ ]);
+ } finally {
+ await previous?.close();
+ await current?.close();
+ files.close();
+ }
+});
+
+async function until(check: () => boolean) {
+ const timeoutMs = 5_000;
+ const deadline = Date.now() + timeoutMs;
+ while (!check()) {
+ if (Date.now() >= deadline) {
+ throw new Error('Released host did not persist the expected state');
+ }
+ await Bun.sleep(1);
+ }
+}
From 16fbe4998fbe8ca7014f9c2eb6aaab76066f3c80 Mon Sep 17 00:00:00 2001
From: massimoalbarello
Date: Fri, 2 Oct 2026 22:35:51 +0100
Subject: [PATCH 5/5] fix: run migrations on the same SQLite connection
---
packages/sync/src/db/client.ts | 13 ++------
packages/sync/src/db/migrate.ts | 39 ++++++++++++++---------
packages/sync/test/cron.test.ts | 28 +++++++++++++++--
packages/sync/test/migration-sql.test.ts | 40 ++++++++++++++++++++++++
4 files changed, 93 insertions(+), 27 deletions(-)
create mode 100644 packages/sync/test/migration-sql.test.ts
diff --git a/packages/sync/src/db/client.ts b/packages/sync/src/db/client.ts
index 8cfbd8d..0ecd75e 100644
--- a/packages/sync/src/db/client.ts
+++ b/packages/sync/src/db/client.ts
@@ -1,18 +1,11 @@
import { Database } from 'bun:sqlite';
-import { DatabaseSync } from 'node:sqlite';
import { runMigrations } from './migrate';
export function openDatabase(path: string): Database {
- // Keep in-memory databases private while both startup connections share the same database.
- const filename =
- path === ':memory:' ? `file:open-sync-${crypto.randomUUID()}?mode=memory&cache=shared` : path;
- // Unlike bun:sqlite in Bun 1.4.0, this driver propagates failures inside SQL batches.
- using migration = new DatabaseSync(filename);
- migration.exec('PRAGMA busy_timeout=5000; PRAGMA journal_mode=WAL;');
- runMigrations(migration);
- const db = new Database(filename, { create: true, strict: true });
+ const db = new Database(path, { create: true, strict: true });
try {
- db.exec('PRAGMA foreign_keys=ON; PRAGMA busy_timeout=5000;');
+ db.exec('PRAGMA foreign_keys=ON; PRAGMA busy_timeout=5000; PRAGMA journal_mode=WAL;');
+ runMigrations(db);
return db;
} catch (error) {
db.close();
diff --git a/packages/sync/src/db/migrate.ts b/packages/sync/src/db/migrate.ts
index f2d0680..a4e624e 100644
--- a/packages/sync/src/db/migrate.ts
+++ b/packages/sync/src/db/migrate.ts
@@ -1,4 +1,4 @@
-import type { DatabaseSync } from 'node:sqlite';
+import type { Database } from 'bun:sqlite';
import initialSchema from './migrations/0000-initial.sql' with { type: 'text' };
import centralizePolling from './migrations/0001-centralize-polling.sql' with { type: 'text' };
@@ -8,30 +8,41 @@ const migrations = [
{ name: '0001-centralize-polling.sql', sql: centralizePolling },
] as const;
-export function runMigrations(db: DatabaseSync): void {
- db.exec('BEGIN IMMEDIATE');
- try {
- db.exec(`CREATE TABLE IF NOT EXISTS __migrations (
+export function runMigrations(db: Database): void {
+ db.transaction(() => {
+ db.query(`CREATE TABLE IF NOT EXISTS __migrations (
name TEXT PRIMARY KEY, applied_at TEXT NOT NULL
- )`);
- const applied = db.prepare('SELECT name FROM __migrations ORDER BY name').all();
+ )`).run();
+ const applied = db
+ .query<{ name: string }, []>('SELECT name FROM __migrations ORDER BY name')
+ .all();
for (const [index, { name }] of applied.entries()) {
if (name !== migrations[index]?.name) {
throw new Error('Unsupported sync database migration history');
}
}
for (const migration of migrations.slice(applied.length)) {
- db.exec(migration.sql);
- db.prepare('INSERT INTO __migrations (name,applied_at) VALUES (?,?)').run(
+ executeSql({ db, sql: migration.sql });
+ db.query('INSERT INTO __migrations (name,applied_at) VALUES (?,?)').run(
migration.name,
new Date().toISOString(),
);
}
- db.exec('COMMIT');
- } catch (error) {
- if (db.isTransaction) {
- db.exec('ROLLBACK');
+ }).immediate();
+}
+
+export function executeSql({ db, sql }: { db: Database; sql: string }): void {
+ // Bun 1.4's exec() can hide failures in non-final statements. Let SQLite parse each
+ // statement instead, preserving trigger bodies, comments and quoted semicolons.
+ // A final no-op also lets prepare() consume files ending in comments or whitespace.
+ let remaining = `${sql}\n;SELECT 1;`;
+ while (remaining.length) {
+ using statement = db.prepare(remaining);
+ const parsed = statement.toString();
+ if (!parsed || !remaining.startsWith(parsed)) {
+ throw new Error('Migration SQL must contain complete statements without parameters');
}
- throw error;
+ statement.run();
+ remaining = remaining.slice(parsed.length);
}
}
diff --git a/packages/sync/test/cron.test.ts b/packages/sync/test/cron.test.ts
index 66035a6..6510855 100644
--- a/packages/sync/test/cron.test.ts
+++ b/packages/sync/test/cron.test.ts
@@ -3,7 +3,7 @@ import { mkdir, stat } from 'node:fs/promises';
import { join } from 'node:path';
import { startCron } from '../src/execution/cron';
import { isolatedScheduler, runRegisteredCron } from './cron-support';
-import { alpha, configure, runtime } from './support';
+import { alpha, configure, fixture, runtime } from './support';
const CRON_FIELDS = 5;
const PRIVATE_DIRECTORY_MODE = 0o700;
@@ -11,13 +11,34 @@ const PRIVATE_SOCKET_MODE = 0o600;
test('a registered shell command invokes the existing headless runtime through its private socket', async () => {
await using table = await isolatedScheduler();
- const f = runtime();
+ const overlapping = Promise.withResolvers();
+ const f = runtime({
+ registration: {
+ ...fixture,
+ load: () => ({
+ async step(context) {
+ await overlapping.promise;
+ return (await fixture.load()).step(context);
+ },
+ }),
+ },
+ });
+ let invocations = 0;
const directory = join(table.directory, 'worker');
const errors: unknown[] = [];
try {
const sync = await configure(f.engine);
await using cron = await startCron({
- runtime: f.engine,
+ runtime: {
+ runDue() {
+ const running = f.engine.runDue();
+ // Both shell processes must reach the host before either poll can finish.
+ if (++invocations === 2) {
+ overlapping.resolve();
+ }
+ return running;
+ },
+ },
directory,
onError: (error) => errors.push(error),
});
@@ -52,6 +73,7 @@ test('a registered shell command invokes the existing headless runtime through i
await cron.close();
expect(await table.table.text()).toBe('');
} finally {
+ overlapping.resolve();
await f.close();
}
});
diff --git a/packages/sync/test/migration-sql.test.ts b/packages/sync/test/migration-sql.test.ts
new file mode 100644
index 0000000..5c9dac0
--- /dev/null
+++ b/packages/sync/test/migration-sql.test.ts
@@ -0,0 +1,40 @@
+import { Database } from 'bun:sqlite';
+import { expect, test } from 'bun:test';
+import { executeSql } from '../src/db/migrate';
+
+test('migration SQL uses SQLite statement boundaries, including triggers, strings and comments', () => {
+ using db = new Database(':memory:');
+ executeSql({
+ db,
+ sql: `
+ -- A comment containing a semicolon ;
+ CREATE TABLE records(value TEXT);
+ CREATE TABLE audit(value TEXT);
+ CREATE TRIGGER record_added AFTER INSERT ON records BEGIN
+ INSERT INTO audit VALUES ('first; value');
+ INSERT INTO audit VALUES (NEW.value);
+ END;
+ /* Another ; comment */
+ INSERT INTO records VALUES ('雪; it''s saved');
+ INSERT INTO records VALUES ('second')
+ -- A trailing comment without a newline`,
+ });
+ expect(db.query('SELECT value FROM audit').all()).toEqual([
+ { value: 'first; value' },
+ { value: "雪; it's saved" },
+ { value: 'first; value' },
+ { value: 'second' },
+ ]);
+});
+
+test('parameterized or invalid migration SQL fails and rolls back prior statements', () => {
+ using db = new Database(':memory:');
+ for (const invalid of ['INSERT INTO records VALUES (?)', 'invalid syntax']) {
+ expect(() =>
+ db
+ .transaction(() => executeSql({ db, sql: `CREATE TABLE records(value TEXT); ${invalid}` }))
+ .immediate(),
+ ).toThrow();
+ expect(db.query('SELECT name FROM sqlite_schema').all()).toEqual([]);
+ }
+});