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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 3 additions & 3 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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),
Expand Down
2 changes: 0 additions & 2 deletions apps/web/backend/src/services/dashboard/service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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.
Expand Down
15 changes: 0 additions & 15 deletions apps/web/frontend/src/routes/_workspace/syncs.$id.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, unknown>,
Expand Down Expand Up @@ -145,10 +144,6 @@ function SyncDetail() {
<dl className="grid grid-cols-[auto_1fr] gap-x-8 gap-y-3 text-sm">
<dt className="text-muted-foreground">Status</dt>
<dd>{sync.status.replaceAll('_', ' ')}</dd>
<dt className="text-muted-foreground">Sync interval</dt>
<dd>{sync.intervalMs / millisecondsPerMinute} minutes</dd>
<dt className="text-muted-foreground">Next scheduled run</dt>
<dd>{nextRun(sync)}</dd>
</dl>
<nav
aria-label="Sync activity"
Expand Down Expand Up @@ -185,16 +180,6 @@ function SyncDetail() {
);
}

function nextRun(sync: { enabled: boolean; status: string; nextDueAt: number }) {
if (!sync.enabled) {
return 'Paused';
}
if (sync.status === 'running') {
return 'After the current run completes';
}
return new Date(sync.nextDueAt).toLocaleString();
}

function syncTitle({
sync,
sourceName,
Expand Down
6 changes: 1 addition & 5 deletions apps/web/frontend/src/routes/_workspace/syncs.index.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -49,11 +49,7 @@ function Syncs() {
<p className="text-muted-foreground text-sm">
{source?.provider && !sync.connection
? 'Waiting for authorization'
: statusLabel(sync.status)}{' '}
·{' '}
{sync.enabled
? `Next poll ${new Date(sync.nextDueAt).toLocaleString()}`
: 'Paused'}
: statusLabel(sync.status)}
</p>
</Link>
</li>
Expand Down
3 changes: 3 additions & 0 deletions bun.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions packages/sync/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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": {
Expand Down
6 changes: 3 additions & 3 deletions packages/sync/src/db/client.ts
Original file line number Diff line number Diff line change
@@ -1,11 +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)).immediate();
db.exec('PRAGMA foreign_keys=ON; PRAGMA busy_timeout=5000; PRAGMA journal_mode=WAL;');
runMigrations(db);
return db;
} catch (error) {
db.close();
Expand Down
48 changes: 48 additions & 0 deletions packages/sync/src/db/migrate.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,48 @@
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' };

// Static imports keep migration history available in both installed packages and compiled hosts.
const migrations = [
{ name: '0000-initial.sql', sql: initialSchema },
{ 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 (
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)) {
executeSql({ db, sql: migration.sql });
db.query('INSERT INTO __migrations (name,applied_at) VALUES (?,?)').run(
migration.name,
new Date().toISOString(),
);
}
}).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');
}
statement.run();
remaining = remaining.slice(parsed.length);
}
}
3 changes: 3 additions & 0 deletions packages/sync/src/db/migrations/0001-centralize-polling.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
ALTER TABLE syncs DROP COLUMN interval_ms;
ALTER TABLE syncs DROP COLUMN next_due_at;
ALTER TABLE syncs DROP COLUMN failure_count;
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
declare const sql: string;
export default sql;
13 changes: 11 additions & 2 deletions packages/sync/src/execution/worker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ export class Worker implements WorkerControl {
#capacityReleased = false;
#closing?: Promise<void>;
#draining?: Promise<void>;
#polling = false;
#drainRequested = false;
#resume?: () => void;
constructor(
Expand Down Expand Up @@ -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<void> {
Expand All @@ -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<void> {
this.ensureOpen();
if (!this.#polling) {
this.input.acquisition.poll();
this.#polling = true;
}
return this.requestDrain();
}
private requestDrain(): Promise<void> {
this.#drainRequested = true;
this.#resume?.();
this.#draining ??= Promise.resolve().then(() => this.drain());
Expand Down Expand Up @@ -98,6 +106,7 @@ export class Worker implements WorkerControl {
} finally {
this.#resume = undefined;
this.#draining = undefined;
this.#polling = false;
}
}
private dispatch(): void {
Expand Down
1 change: 0 additions & 1 deletion packages/sync/src/http/controller.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 }) };
Expand Down
2 changes: 1 addition & 1 deletion packages/sync/src/models/definition.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
7 changes: 0 additions & 7 deletions packages/sync/src/models/sync.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,8 +11,6 @@ export interface Sync {
destination: { type: string; config: JsonObject };
enabled: boolean;
checkpoint: JsonValue;
intervalMs: number;
nextDueAt: number;
status: SyncStatus;
errorCode: string | null;
}
Expand All @@ -21,7 +19,6 @@ export interface CreateSync extends Scope {
connection?: ConnectionRef;
config: JsonObject;
destination: { type: string; input: JsonObject };
intervalMs?: number;
enabled?: boolean;
}

Expand Down Expand Up @@ -51,8 +48,6 @@ export interface SyncSummary {
destinationType: string;
connection?: ConnectionRef;
enabled: boolean;
intervalMs: number;
nextDueAt: number;
status: SyncStatus;
errorCode: string | null;
}
Expand All @@ -63,8 +58,6 @@ export function summarizeSync(sync: Sync): SyncSummary {
destinationType: sync.destination.type,
connection: sync.connection,
enabled: sync.enabled,
intervalMs: sync.intervalMs,
nextDueAt: sync.nextDueAt,
status: sync.status,
errorCode: sync.errorCode,
};
Expand Down
4 changes: 1 addition & 3 deletions packages/sync/src/repositories/acquisition/contract.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,9 +6,9 @@ export interface AcquisitionLease extends Scope {
sync: Sync;
generation: number;
force: boolean;
failureCount: number;
}
export interface AcquisitionRepository {
poll(): void;
capacityReleased(): void;
claim(leaseMs: number): AcquisitionLease | undefined;
hasCapacity(): boolean;
Expand All @@ -17,8 +17,6 @@ export interface AcquisitionRepository {
lease: AcquisitionLease;
state: Exclude<SyncStatus, 'running' | 'disabled'>;
errorCode?: string;
delay: number;
failureCount?: number;
pause?: boolean;
}): void;
}
22 changes: 9 additions & 13 deletions packages/sync/src/repositories/acquisition/lease.ts
Original file line number Diff line number Diff line change
Expand Up @@ -30,22 +30,21 @@ 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 status='ready',error_code='lease_expired',expires_at=NULL
WHERE enabled=1 AND expires_at<=?`).run(now);
const row = db
.query<
{
id: string;
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 expires_at IS NULL AND next_due_at<=?
ORDER BY next_due_at,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;
}
Expand All @@ -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();
Expand All @@ -88,15 +86,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=?,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,
now + input.delay,
Number(input.pause ?? false),
input.failureCount ?? null,
Number(input.state === 'succeeded'),
lease.ownerId,
lease.sync.id,
Expand Down
13 changes: 9 additions & 4 deletions packages/sync/src/repositories/acquisition/sqlite.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,10 +15,17 @@ 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 IN ('succeeded','retrying','interrupted','waiting_for_capacity')",
)
.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 status='ready' WHERE enabled=1 AND status='waiting_for_capacity'")
.run();
}
claim(leaseMs: number) {
return claimAcquisition({ ...this.input, leaseMs });
Expand Down Expand Up @@ -52,8 +59,6 @@ export class SqliteAcquisition implements AcquisitionRepository {
db,
lease: input.lease,
state: page.complete ? 'succeeded' : 'ready',
delay: page.complete ? sync.intervalMs : 0,
failureCount: page.complete ? 0 : undefined,
});
}).immediate();
}
Expand Down
4 changes: 2 additions & 2 deletions packages/sync/src/repositories/assets/sqlite.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 status='ready' 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[] {
Expand Down
Loading
Loading