diff --git a/apps/cli/README.md b/apps/cli/README.md index fa9f0518..bdd8fa97 100644 --- a/apps/cli/README.md +++ b/apps/cli/README.md @@ -12,11 +12,16 @@ Requires Node.js >= 22.5 (`node:sqlite`, global `fetch`/`WebSocket`/WebCrypto). synch login # device-code sign-in (prints URL + code) synch logout # sign out, clear stored keys synch vault connect --vault-id # unlock a remote vault for a directory +synch pull # download only; never upload local changes synch sync # one-shot synchronization synch watch # keep syncing until interrupted synch status # account, vault, and sync state ``` +`synch pull` never scans for local changes and never uploads pending local +mutations. Remote versions replace differing files in the target directory, so +use it only for read-only replicas or backup staging directories. + Common options: `--vault ` (default: current directory) and `--api-url ` (or the `SYNCH_API_URL` environment variable). diff --git a/apps/cli/src/commands/pull.ts b/apps/cli/src/commands/pull.ts new file mode 100644 index 00000000..dc3a5c2b --- /dev/null +++ b/apps/cli/src/commands/pull.ts @@ -0,0 +1,14 @@ +import type { CliAppContext } from "../app/context"; + +export async function runPull(ctx: CliAppContext): Promise { + await ctx.initializeAuth(); + ctx.requireVerifiedAuth(); + await ctx.openVaultSession(); + + ctx.logger.log(`Pulling remote changes into ${ctx.vaultPath} ...`); + await ctx.engine.pullOnlyOnce(); + ctx.logger.log( + `Pull complete (${ctx.syncProgress.completedEntries}/${ctx.syncProgress.totalEntries} entries).`, + ); + return 0; +} diff --git a/apps/cli/src/main.ts b/apps/cli/src/main.ts index 4c6454e0..3c7030fd 100644 --- a/apps/cli/src/main.ts +++ b/apps/cli/src/main.ts @@ -3,6 +3,7 @@ import { parseArgs } from "node:util"; import { CliAppContext, CliUsageError, describeError } from "./app/context"; import { runLogin } from "./commands/login"; import { runLogout } from "./commands/logout"; +import { runPull } from "./commands/pull"; import { runStatus } from "./commands/status"; import { runSync } from "./commands/sync"; import { runVaultConnect } from "./commands/vault-connect"; @@ -16,6 +17,7 @@ Usage: synch login Sign in with a device code synch logout Sign out and clear stored keys synch vault connect --vault-id Connect a vault directory to a remote vault + synch pull Download remote changes without uploading local changes synch sync Synchronize the vault once and exit synch watch Keep the vault in sync until interrupted synch status Show account, vault, and sync state @@ -55,6 +57,8 @@ async function main(argv: string[]): Promise { return await runLogout(ctx); case "vault-connect": return await runVaultConnect(ctx, values["vault-id"]); + case "pull": + return await runPull(ctx); case "sync": return await runSync(ctx); case "watch": @@ -99,6 +103,7 @@ interface CliParseArgsConfig { type CliCommand = | "login" | "logout" + | "pull" | "vault-connect" | "sync" | "watch" @@ -109,6 +114,7 @@ function resolveCommand(positionals: string[]): CliCommand | null { switch (first) { case "login": case "logout": + case "pull": case "sync": case "watch": case "status": diff --git a/packages/sync-client/src/sync/engine/__tests__/auto-sync/pull-only.test.ts b/packages/sync-client/src/sync/engine/__tests__/auto-sync/pull-only.test.ts new file mode 100644 index 00000000..795e6936 --- /dev/null +++ b/packages/sync-client/src/sync/engine/__tests__/auto-sync/pull-only.test.ts @@ -0,0 +1,143 @@ +import { describe, expect, it, vi } from "vitest"; + +import { createTestSyncStore } from "../../../../test-support/in-memory-sync-store"; +import { SyncAutoLoop } from "../../auto-sync"; +import { + createPushResult, + createRealtimeClient, + createToken, +} from "./helpers"; + +describe("SyncAutoLoop pull-only", () => { + it("pulls once without scheduling pending local mutations for push", async () => { + const store = createTestSyncStore(); + await store.setCursor(4); + const pullOnce = vi.fn(async () => {}); + const pushPendingMutations = vi.fn(async () => createPushResult()); + let sessionClosed = false; + const realtimeClient = createRealtimeClient( + undefined, + (session) => { + session.close = () => { + sessionClosed = true; + }; + }, + 7, + ); + const openSession = vi.fn(realtimeClient.openSession); + const autoLoop = new SyncAutoLoop({ + getApiBaseUrl: () => "http://127.0.0.1:8787", + getSyncToken: async () => createToken(), + getSyncStore: () => store, + pushPendingMutations, + pullOnce, + realtimeClient: { openSession }, + }); + + await autoLoop.pullOnlyOnce(); + + expect(openSession).toHaveBeenCalledWith( + "http://127.0.0.1:8787", + expect.objectContaining({ vaultId: "vault-1" }), + 4, + expect.any(Object), + ); + expect(pullOnce).toHaveBeenCalledTimes(1); + expect(pushPendingMutations).not.toHaveBeenCalled(); + expect(sessionClosed).toBe(true); + await store.close(); + }); + + it("refuses to run alongside the auto-sync loop", async () => { + const store = createTestSyncStore(); + const autoLoop = new SyncAutoLoop({ + getApiBaseUrl: () => "http://127.0.0.1:8787", + getSyncToken: async () => createToken(), + getSyncStore: () => store, + pushPendingMutations: vi.fn(async () => createPushResult()), + pullOnce: vi.fn(async () => {}), + realtimeClient: createRealtimeClient(), + }); + + await autoLoop.start(); + await expect(autoLoop.pullOnlyOnce()).rejects.toThrow( + "requires the auto-sync loop and all in-flight sync work to be stopped", + ); + autoLoop.stop(); + await store.close(); + }); + + it("refuses to pull when the store contains pending local changes", async () => { + const store = createTestSyncStore(); + await store.markEntryDirty({ + mutationId: "mutation-1", + entryId: "entry-1", + op: "delete", + baseRevision: 0, + blobId: null, + hash: null, + encryptedMetadata: "encrypted-metadata", + createdAt: 1, + }); + const pullOnce = vi.fn(async () => {}); + const openSession = vi.fn(createRealtimeClient().openSession); + const autoLoop = new SyncAutoLoop({ + getApiBaseUrl: () => "http://127.0.0.1:8787", + getSyncToken: async () => createToken(), + getSyncStore: () => store, + pushPendingMutations: vi.fn(async () => createPushResult()), + pullOnce, + realtimeClient: { openSession }, + }); + + await expect(autoLoop.pullOnlyOnce()).rejects.toThrow( + "no pending local changes", + ); + expect(openSession).not.toHaveBeenCalled(); + expect(pullOnce).not.toHaveBeenCalled(); + await store.close(); + }); + + it("explains how to rebuild state when the local cursor is ahead", async () => { + const store = createTestSyncStore(); + await store.setCursor(8); + const autoLoop = new SyncAutoLoop({ + getApiBaseUrl: () => "http://127.0.0.1:8787", + getSyncToken: async () => createToken(), + getSyncStore: () => store, + pushPendingMutations: vi.fn(async () => createPushResult()), + pullOnce: vi.fn(async () => {}), + realtimeClient: createRealtimeClient(undefined, undefined, 7), + }); + + await expect(autoLoop.pullOnlyOnce()).rejects.toThrow( + "Move .synch/sync.sqlite aside", + ); + await store.close(); + }); + + it("propagates asynchronous session errors and still closes the session", async () => { + const store = createTestSyncStore(); + let sessionClosed = false; + const sessionError = new Error("session failed"); + const autoLoop = new SyncAutoLoop({ + getApiBaseUrl: () => "http://127.0.0.1:8787", + getSyncToken: async () => createToken(), + getSyncStore: () => store, + pushPendingMutations: vi.fn(async () => createPushResult()), + pullOnce: vi.fn(async () => {}), + realtimeClient: createRealtimeClient( + (callbacks) => callbacks.onError(sessionError), + (session) => { + session.close = () => { + sessionClosed = true; + }; + }, + ), + }); + + await expect(autoLoop.pullOnlyOnce()).rejects.toThrow("session failed"); + expect(sessionClosed).toBe(true); + await store.close(); + }); +}); diff --git a/packages/sync-client/src/sync/engine/auto-sync.ts b/packages/sync-client/src/sync/engine/auto-sync.ts index 5273f4cc..d4ef024a 100644 --- a/packages/sync-client/src/sync/engine/auto-sync.ts +++ b/packages/sync-client/src/sync/engine/auto-sync.ts @@ -14,7 +14,7 @@ import { type SyncStorageStatus, type PresenceUpdatedPush, } from "../remote/realtime-client"; -import type { SyncCursorStore } from "../store/ports"; +import type { SyncCursorStore, SyncMutationStore } from "../store/ports"; import { SyncAutoLoopState, type SyncConnectionState } from "./auto-sync-state"; import { AutoSyncTimers } from "./auto-sync-timers"; import { PendingSyncWorkQueue } from "./auto-sync-work-queue"; @@ -28,7 +28,7 @@ const DEFAULT_SYNC_RETRY_MAX_DELAY_MS = 30_000; export interface SyncAutoLoopDeps { getApiBaseUrl: () => string; getSyncToken: () => Promise; - getSyncStore: () => SyncCursorStore | null; + getSyncStore: () => (SyncCursorStore & SyncMutationStore) | null; pushPendingMutations: ( session: SyncRealtimeSession, shouldYield: () => boolean, @@ -163,6 +163,72 @@ export class SyncAutoLoop { ); } + /** + * Pull remote changes once without reconciling or uploading local changes. + * + * This deliberately uses a short-lived realtime session outside the normal + * auto-sync drain loop, so pending local mutations are never scheduled for + * push. It is intended for read-only replicas such as backup hosts. + */ + async pullOnlyOnce(): Promise { + if (this.isActive() || this.connectPromise || this.drainPromise) { + throw new Error( + "Pull-only sync requires the auto-sync loop and all in-flight sync work to be stopped.", + ); + } + + const store = this.deps.getSyncStore(); + if (!store) { + throw new Error("Sync store is not initialized."); + } + + const pendingMutations = await store.listDirtyEntries(1); + if (pendingMutations.length > 0) { + throw new Error( + "Pull-only sync requires a read-only replica with no pending local changes. Use a fresh backup directory or resolve the pending changes with `synch sync` first.", + ); + } + + const token = await this.deps.getSyncToken(); + const cursor = await store.getCursor(); + let sessionError: Error | null = null; + const session = await this.realtimeClient.openSession( + this.deps.getApiBaseUrl(), + token, + cursor, + { + onCursorAdvanced() {}, + onStorageStatusUpdated() {}, + onPolicyUpdated() {}, + onPresenceUpdated() {}, + onPresenceCleared() {}, + onPresenceAvailabilityChanged() {}, + onClose() {}, + onError(error) { + sessionError ??= error; + }, + }, + ); + + try { + if (cursor > session.serverCursor) { + throw new SyncRealtimeError( + "cursor_ahead_of_server", + "This device's sync history no longer matches the remote vault. Move .synch/sync.sqlite aside, then run `synch vault connect --vault-id ` to rebuild this read-only replica's sync state.", + ); + } + if (sessionError) { + throw sessionError; + } + await this.deps.pullOnce(session); + if (sessionError) { + throw sessionError; + } + } finally { + session.close(); + } + } + requestPull(targetCursor: number | null = null): void { if (!this.isActive()) { return; diff --git a/packages/sync-client/src/sync/runtime/sync-engine.ts b/packages/sync-client/src/sync/runtime/sync-engine.ts index d9193336..1fc0c4c5 100644 --- a/packages/sync-client/src/sync/runtime/sync-engine.ts +++ b/packages/sync-client/src/sync/runtime/sync-engine.ts @@ -489,6 +489,10 @@ export class SyncEngine { return await this.syncAutoLoop.syncNow(); } + async pullOnlyOnce(): Promise { + await this.syncAutoLoop.pullOnlyOnce(); + } + async flushDebouncedPushAndWaitForInFlight(): Promise { await this.waitForLocalMutationWork(); this.syncAutoLoop.flushDebouncedPush();