Skip to content
Merged
24 changes: 24 additions & 0 deletions .changeset/t12996-cloud-sync.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,24 @@
---
id: t12996-cloud-sync
tasks: [T12996]
kind: feat
summary: "`cleo cloud sync` - seal, push, pull and apply each attached stream in one call, one result per stream"
---

`cleo cloud sync [--scope project|global]` runs the change journal end to end for each attached stream: the project's and the
account's global store, or just the one `--scope` names (T12996).

- **Per stream:** it seals pending writes, pushes them as segments (`pushSyncStream`), pulls the stream's new segments, and applies
them (`pullSyncStream`).
- **One LAFS envelope** (`CloudSyncResult`) reports, per stream:
- `status`: `synced`, `disabled`, `paused` (device clock ahead), `not-attached`, `refused`, or `failed` (an error stopped this
stream, reported, never thrown, so the next stream still syncs);
- `skipped`: legs with nothing to do (a flag off, or push before this store has its own genesis), which never count as refused;
- sealed, built, sent, duplicates, received, staged, redelivered, applied, held and conflict counts;
- the staged position against the server's head.
- **`E_SYNC_DISABLED`** is returned, naming `cleo sync enable push`, when no attached stream has `sync.push` or `sync.pull` on. A store
with `sync.pull` off now refuses its pull before it looks for a pull position.
- **An interrupted run resumes on the next:** a segment is resent with the same bytes, and a transaction is never staged or applied
twice.
- **One connection and one account-key unlock serve every stream of a run.**
- `sync.push` and `sync.pull` stay unreleased: the CLI never opts in, so stores can't turn them on yet.
21 changes: 21 additions & 0 deletions packages/cleo/src/cli/commands/cloud.ts
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@ import {
runCloudPull,
runCloudPush,
runCloudRestore,
runCloudSync,
runCloudVault,
runCloudVerify,
} from '../lib/nexus-vault-cli.js';
Expand Down Expand Up @@ -128,6 +129,25 @@ const projectsSubCommand = defineCommand({
},
});

const syncSubCommand = defineCommand({
meta: {
name: 'sync',
description:
"Sync the change journal with Cleo Nexus: seal pending writes, push them as segments, pull the stream's new segments and apply them, for each attached stream (the project's and the account's global store; --scope picks one). One envelope with per-stream sealed, sent, received, applied, held and conflict counts and the staged position against the server's head. Refused with E_SYNC_DISABLED until `cleo sync enable push`. An interrupted run resumes on the next: nothing is sent or applied twice.",
},
args: {
scope: {
type: 'string',
description: "One stream: 'project' or 'global' (default: every attached stream).",
},
'api-url': NEXUS_API_URL_ARG,
json: JSON_ARG,
},
async run({ args }) {
await runCloudSync(args as Record<string, unknown>);
},
});

const pushSubCommand = defineCommand({
meta: {
name: 'push',
Expand Down Expand Up @@ -320,6 +340,7 @@ export const cloudCommand = defineCommand({
projects: projectsSubCommand,
activity: activitySubCommand,
conflicts: conflictsSubCommand,
sync: syncSubCommand,
push: pushSubCommand,
pull: pullSubCommand,
restore: restoreSubCommand,
Expand Down
35 changes: 35 additions & 0 deletions packages/cleo/src/cli/lib/__tests__/nexus-vault-cli.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest';

const pushNexusVault = vi.fn();
const enableSyncPush = vi.fn();
const cloudSync = vi.fn();
const restoreNexusVault = vi.fn();
const verifyNexusVault = vi.fn();
const nexusVaultStatus = vi.fn();
Expand All @@ -23,6 +24,7 @@ const resolveNexusProjectRef = vi.fn();
const assertNexusRestoreTarget = vi.fn();

vi.mock('@cleocode/core/cloud/nexus-vault.js', () => ({
cloudSync,
enableSyncPush,
pushNexusVault,
restoreNexusVault,
Expand All @@ -49,6 +51,7 @@ const {
runCloudPull,
runCloudPush,
runCloudRestore,
runCloudSync,
runCloudVault,
runCloudVerify,
runSyncEnablePush,
Expand Down Expand Up @@ -181,6 +184,38 @@ describe('flags reach the core calls', () => {
expect(written()).toContain('cp-1');
});

it('cloud sync passes the API URL, and a scope only when one is given (T12996)', async () => {
cloudSync.mockReset().mockResolvedValue({
apiUrl: API,
streams: [
{
scope: 'project',
streamId: 'project:p',
status: 'synced',
refused: null,
sealed: 1,
built: 1,
sent: 1,
duplicates: 0,
received: 2,
staged: 2,
redelivered: 0,
applied: 2,
held: 0,
conflicts: 0,
after: 7,
head: 7,
},
],
warnings: [],
});
await runCloudSync({ 'api-url': API });
expect(opts(cloudSync)).toEqual({ apiUrl: API });
await runCloudSync({ scope: 'global' });
expect(cloudSync.mock.calls[1]?.[0]).toEqual({ apiUrl: undefined, scope: 'global' });
expect(written()).toContain('project:p');
});

it('restore passes checkpoint, project, into and force', async () => {
await runCloudRestore({
scope: 'project',
Expand Down
27 changes: 27 additions & 0 deletions packages/cleo/src/cli/lib/nexus-vault-cli.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ import type {
CloudPushResult,
CloudRestoreResult,
CloudSyncPushEnableResult,
CloudSyncResult,
CloudVaultScope,
CloudVaultStatusResult,
CloudVerifyResult,
Expand Down Expand Up @@ -108,6 +109,32 @@ export async function runSyncEnablePush(args: Args): Promise<void> {
);
}

/**
* `cleo cloud sync [--scope]`: seal, push, pull and apply each attached
* stream (T12996). Without `--scope`, every attached stream.
*
* @param args - Parsed args.
*/
export async function runCloudSync(args: Args): Promise<void> {
const scope = stringArg(args, 'scope') === undefined ? undefined : scopeArg(args, 'cloud.sync');
await runCloudRead<CloudSyncResult>(
'cloud.sync',
async () =>
(await vaultModule()).cloudSync({
apiUrl: nexusApiUrlArg(args),
...(scope !== undefined ? { scope } : {}),
}),
(r) =>
r.streams
.map((st) =>
st.status === 'synced'
? `${st.streamId}: sent ${st.sent} segment(s), received ${st.received}, applied ${st.applied}${st.held > 0 ? `, ${st.held} held` : ''}${st.conflicts > 0 ? `, ${st.conflicts} in conflict` : ''} (at ${st.after} of ${st.head}).`
: `${st.streamId ?? st.scope}: ${st.status}${st.refused ? ` (${st.refused})` : ''}.`,
)
.join('\n'),
);
}

/**
* `cleo cloud pull [--scope] [--force]`.
*
Expand Down
2 changes: 2 additions & 0 deletions packages/contracts/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3157,6 +3157,8 @@ export {
type CloudPushResult,
type CloudRestoreResult,
type CloudSyncPushEnableResult,
type CloudSyncResult,
type CloudSyncStreamResult,
type CloudVaultLease,
type CloudVaultScope,
type CloudVaultSnapshot,
Expand Down
2 changes: 2 additions & 0 deletions packages/contracts/src/nexus-account.ts
Original file line number Diff line number Diff line change
Expand Up @@ -109,6 +109,8 @@ export const NEXUS_ACCOUNT_ERROR_CODES = [
* checkpoint, then joins the stream (T12999), never writing a second genesis.
*/
'E_NEXUS_SYNC_STREAM_JOURNALED',
/** `cleo cloud sync` found no stream with `sync.push` or `sync.pull` on (`cleo sync enable`). */
'E_SYNC_DISABLED',
// Projects by name (`cleo cloud restore <name>`, T13102).
/** No project of the account has that name, label or id. */
'E_NEXUS_PROJECT_NOT_FOUND',
Expand Down
43 changes: 43 additions & 0 deletions packages/contracts/src/nexus-vault.ts
Original file line number Diff line number Diff line change
Expand Up @@ -236,6 +236,49 @@ export interface CloudSyncPushEnableResult {
warnings: CloudWarning[];
}

/** One stream of `cleo cloud sync` (T12996). */
export interface CloudSyncStreamResult {
scope: CloudVaultScope;
streamId: string | null;
/**
* `synced`: pushed and pulled what was enabled; `disabled`: neither `sync.push` nor `sync.pull` is on;
* `paused`: push paused for a device clock ahead of the server's; `not-attached`: this machine has no
* store on that stream; `refused`: a precondition or the server refused (see `refused`); `failed`: an
* error stopped this stream (see `refused`); the other streams still synced.
*/
status: 'synced' | 'disabled' | 'paused' | 'not-attached' | 'refused' | 'failed';
/** Why this stream did not sync fully, or null. */
refused: string | null;
/**
* Legs with nothing to do, not failures: a flag off, or push before this store has its own genesis
* (a store that only pulls).
*/
skipped: string[];
/** Transactions sealed, segments persisted, segments the server stored (`duplicates` of them retries). */
sealed: number;
built: number;
sent: number;
duplicates: number;
/** Segments received, transactions staged, re-deliveries skipped. */
received: number;
staged: number;
redelivered: number;
/** Transactions applied; still held (pending, clock skew, or a newer schema); in conflict. */
applied: number;
held: number;
conflicts: number;
/** The stream position this store has staged up to, and the server's head (null when not pulled). */
after: number | null;
head: number | null;
}

/** `cleo cloud sync` (T12996): seal, push, pull and apply each attached stream. */
export interface CloudSyncResult {
apiUrl: string;
streams: CloudSyncStreamResult[];
warnings: CloudWarning[];
}

/** `cleo cloud pull` / `cleo cloud restore`. */
export interface CloudRestoreResult {
apiUrl: string;
Expand Down
83 changes: 83 additions & 0 deletions packages/core/src/cloud/__tests__/nexus-vault.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -121,6 +121,7 @@ import { hasUnsyncedNexusBackup, runNexusFirstRun } from '../nexus-first-run.js'
import { linkProjectToNexus } from '../nexus-link.js';
import { listNexusNamedProjects, resolveNexusProjectRef } from '../nexus-project-names.js';
import {
cloudSync,
enableSyncPush,
nexusVaultStatus,
pullSyncStream,
Expand Down Expand Up @@ -5536,6 +5537,88 @@ describe('sync enable push: the genesis checkpoint (T12343 S4-1b)', () => {
expect(again).toMatchObject({ segments: 0, staged: 0 });
});

it('C-1: cloud sync seals, pushes, pulls and applies the stream in one call; a rerun does nothing', async () => {
const { m, dbPath } = await journalMachine();
await on(m, () => enableSyncPush(vopts(m, { allowUnreleased: true })));
await on(m, async () => {
const db = await storeOf(dbPath);
setSyncFlag(db, 'sync.pull', true, { allowUnreleased: true });
db.exec(
"INSERT INTO tasks_tasks (id, title, type, status, priority, uid, birth_fp) VALUES ('C1', 'title C1', 'task', 'pending', 'medium', 'uid-C1', 'fp-C1')",
);
});
const r = await on(m, () => cloudSync(vopts(m, { scope: 'project', allowUnreleased: true })));
expect(r.streams).toHaveLength(1);
expect(r.streams[0]).toMatchObject({
scope: 'project',
streamId: STREAM,
status: 'synced',
refused: null,
sealed: 1,
sent: 1,
received: 1,
staged: 1,
applied: 1,
held: 0,
conflicts: 0,
});
expect(r.streams[0]?.after).toBe(r.streams[0]?.head);
const again = await on(m, () =>
cloudSync(vopts(m, { scope: 'project', allowUnreleased: true })),
);
expect(again.streams[0]).toMatchObject({ status: 'synced', sealed: 0, sent: 0, received: 0 });
});

it('C-1: an interrupted sync resumes on the next run with nothing sent or applied twice', async () => {
const { m, dbPath } = await journalMachine();
await on(m, () => enableSyncPush(vopts(m, { allowUnreleased: true })));
await on(m, async () => {
const db = await storeOf(dbPath);
setSyncFlag(db, 'sync.pull', true, { allowUnreleased: true });
db.exec(
"INSERT INTO tasks_tasks (id, title, type, status, priority, uid, birth_fp) VALUES ('C2', 'title C2', 'task', 'pending', 'medium', 'uid-C2', 'fp-C2')",
);
});
fake.beforeSegment = {
deviceId: DEVICE_A,
run: async () => {
throw new ApiFail(403, 'E_FORBIDDEN');
},
};
// The failure is this stream's, reported, never thrown: every other stream still syncs.
const failed = await on(m, () => cloudSync(vopts(m, { allowUnreleased: true })));
expect(failed.streams.map((x) => x.scope)).toEqual(['project', 'global']);
expect(failed.streams[0]).toMatchObject({ status: 'failed', sent: 0 });
expect(failed.streams[0]?.refused).toMatch(/E_FORBIDDEN/);
const before = fake.stream(STREAM).segments.length;
const r = await on(m, () => cloudSync(vopts(m, { scope: 'project', allowUnreleased: true })));
expect(r.streams[0]).toMatchObject({
status: 'synced',
built: 0,
sent: 1,
received: 1,
applied: 1,
});
expect(fake.stream(STREAM).segments).toHaveLength(before + 1);
await on(m, async () => {
const db = await storeOf(dbPath);
expect((db.prepare('SELECT count(*) AS n FROM _sync_inbox').get() as { n: number }).n).toBe(
1,
);
expect(
(db.prepare("SELECT count(*) AS n FROM tasks_tasks WHERE id = 'C2'").get() as { n: number })
.n,
).toBe(1);
});
});

it('C-1: with the journal flags off, cloud sync refuses with E_SYNC_DISABLED naming the remedy', async () => {
const { m } = await journalMachine();
const refused = await failure(on(m, () => cloudSync(vopts(m, { scope: 'project' }))));
expect(refused.code).toBe('E_SYNC_DISABLED');
expect(refused.fix).toContain('cleo sync enable push');
});

it('S5-1: a store that synced a vault snapshot from before the journal genesis refuses to pull across it (T13306)', async () => {
const { m: a, dbPath } = await journalMachine();
await on(a, () => pushNexusVault(vopts(a))); // a vault snapshot, before any journal
Expand Down
Loading
Loading