diff --git a/package.json b/package.json index df277592..e9f69c34 100644 --- a/package.json +++ b/package.json @@ -86,6 +86,7 @@ "test:e2e:obsidian:security-seed-reconnect": "tsx test/e2e-obsidian/scripts/security-seed-reconnect.ts", "test:e2e:obsidian:hidden-file-snippet-sync": "tsx test/e2e-obsidian/scripts/hidden-file-snippet-sync.ts", "test:e2e:obsidian:customisation-sync": "tsx test/e2e-obsidian/scripts/customisation-sync.ts", + "test:e2e:obsidian:received-change-readiness": "tsx test/e2e-obsidian/scripts/received-change-readiness.ts", "test:e2e:obsidian:remote-feature-change": "tsx test/e2e-obsidian/scripts/remote-feature-change.ts", "test:e2e:obsidian:internal-metadata-migration": "tsx test/e2e-obsidian/scripts/internal-metadata-migration.ts", "test:e2e:obsidian:setting-markdown-export": "tsx test/e2e-obsidian/scripts/setting-markdown-export.ts", diff --git a/src/serviceFeatures/replication/ReplicateResultProcessor.ts b/src/serviceFeatures/replication/ReplicateResultProcessor.ts index f0010c96..9474e58e 100644 --- a/src/serviceFeatures/replication/ReplicateResultProcessor.ts +++ b/src/serviceFeatures/replication/ReplicateResultProcessor.ts @@ -108,8 +108,15 @@ export class ReplicateResultProcessor { } public resume() { this._suspended = false; + this.continueHeldDocuments(); + } + /** + * Continue the queued documents which were held, for example while the application was not ready. + * An explicit suspension, by `suspend()` or by the settings, remains in effect. + */ + public continueHeldDocuments() { this.updateProcessingActivity(); - fireAndForget(() => this.runProcessQueue()); + this.triggerProcessQueue(); } // Whether the processing is suspended diff --git a/src/serviceFeatures/replication/ReplicateResultProcessor.unit.spec.ts b/src/serviceFeatures/replication/ReplicateResultProcessor.unit.spec.ts index 2ea388ed..4e992692 100644 --- a/src/serviceFeatures/replication/ReplicateResultProcessor.unit.spec.ts +++ b/src/serviceFeatures/replication/ReplicateResultProcessor.unit.spec.ts @@ -207,6 +207,35 @@ describe("ReplicateResultProcessor", () => { expect(isReady).toHaveBeenCalledOnce(); }); + it("applies documents held before readiness once the application becomes ready", async () => { + const { isReady, processor, processSynchroniseResult } = setup({ applicationReady: false }); + processor.enqueueAll([note("held")]); + await new Promise((resolve) => setTimeout(resolve, 0)); + expect(processSynchroniseResult).not.toHaveBeenCalled(); + expect(processor["_queuedChanges"]).toHaveLength(1); + + isReady.mockReturnValue(true); + processor.continueHeldDocuments(); + + await vi.waitFor(() => expect(processSynchroniseResult).toHaveBeenCalledOnce()); + expect(processor["_queuedChanges"]).toHaveLength(0); + }); + + it("keeps an explicit suspension when the application becomes ready", async () => { + const { isReady, processor, processSynchroniseResult } = setup({ applicationReady: false }); + processor.suspend(); + processor.enqueueAll([note("held")]); + + isReady.mockReturnValue(true); + processor.continueHeldDocuments(); + await new Promise((resolve) => setTimeout(resolve, 0)); + expect(processSynchroniseResult).not.toHaveBeenCalled(); + expect(processor["_queuedChanges"]).toHaveLength(1); + + processor.resume(); + await vi.waitFor(() => expect(processSynchroniseResult).toHaveBeenCalledOnce()); + }); + it("applies results in remediation mode, which never reports readiness", () => { const { processor } = setup({ applicationReady: false, diff --git a/src/serviceFeatures/replication/index.ts b/src/serviceFeatures/replication/index.ts index 74a4f1e4..37eb5f52 100644 --- a/src/serviceFeatures/replication/index.ts +++ b/src/serviceFeatures/replication/index.ts @@ -4,6 +4,7 @@ import { UnresolvedErrorManager } from "@vrtmrz/livesync-commonlib/compat/servic import type { ServiceContext } from "@vrtmrz/livesync-commonlib/context"; import { fireAndForget } from "octagonal-wheels/promises"; import type { IMinimumLiveSyncCommands, LiveSyncBaseCore } from "@/LiveSyncBaseCore"; +import { EVENT_APPLICATION_READY } from "@/common/events"; import { createAutomaticReplicationTriggers } from "./automaticTriggers"; import { createCentralCompatibilityRecovery } from "./centralCompatibilityRecovery"; import { createOnlineReplicationPreflight, createSecuritySeedPreflight } from "./preflight"; @@ -102,6 +103,9 @@ export function useReplicationFeature resultProcessor.restoreFromSnapshotOnce()); return Promise.resolve(true); }); + // Commonlib emits this each time it establishes readiness. Documents held until then, such as those restored + // from the snapshot or received during a fetch, continue from here. + services.context.events.onEvent(EVENT_APPLICATION_READY, () => resultProcessor.continueHeldDocuments()); services.appLifecycle.onSettingLoaded.addHandler(initialiseAutomaticReplicationTriggers); services.replication.parseSynchroniseResult.addHandler((documents) => { resultProcessor.enqueueAll(documents); diff --git a/src/serviceFeatures/replication/receivedChangeReadiness.unit.spec.ts b/src/serviceFeatures/replication/receivedChangeReadiness.unit.spec.ts new file mode 100644 index 00000000..1ebf14f8 --- /dev/null +++ b/src/serviceFeatures/replication/receivedChangeReadiness.unit.spec.ts @@ -0,0 +1,171 @@ +import { createServiceContext } from "@vrtmrz/livesync-commonlib/context"; +import { VERSIONING_DOCID, type EntryDoc } from "@vrtmrz/livesync-commonlib/compat/common/types"; +import { EVENT_SETTING_SAVED } from "@vrtmrz/livesync-commonlib/compat/events/coreEvents"; +import { describe, expect, it, vi } from "vitest"; +import { EVENT_APPLICATION_READY, eventHub } from "@/common/events"; +import { useReplicationFeature } from "./index"; + +type ParseHandler = (documents: PouchDB.Core.ExistingDocument[]) => Promise; + +function receivedNote(id: string): PouchDB.Core.ExistingDocument { + return { + _id: id, + _rev: "1-received", + path: `${id}.md`, + ctime: 1, + mtime: 2, + size: 1, + children: [], + datatype: "plain", + type: "plain", + eden: {}, + } as unknown as PouchDB.Core.ExistingDocument; +} + +function setup() { + let applicationReady = false; + const settings = { + handleFilenameCaseSensitive: false, + ignoreFiles: "", + maxMTimeForReflectEvents: 0, + suspendParseReplicationResult: false, + syncIgnoreRegEx: "", + syncInternalFiles: false, + syncMaxSizeInMB: 0, + syncOnlyRegEx: "", + useIgnoreFiles: false, + }; + const processSynchroniseResult = vi.fn(async () => true); + const settingLoadedHandlers: (() => Promise)[] = []; + const context = createServiceContext(); + const keyValueDB = { + get: vi.fn(async () => undefined), + set: vi.fn(async () => undefined), + }; + const localDatabase = { + getRaw: vi.fn(async (id: string) => { + if (id === VERSIONING_DOCID) throw { status: 404 }; + return { _id: id, _rev: "1-received" }; + }), + getDBEntryFromMeta: vi.fn(async (entry: object) => ({ ...entry, data: "received content" })), + }; + let parseHandler: ParseHandler | undefined; + const services = { + API: { isMobile: vi.fn(() => false), isOnline: true }, + appLifecycle: { + getUnresolvedMessages: { addHandler: vi.fn() }, + isReady: () => applicationReady, + isSuspended: vi.fn(() => false), + onSettingLoaded: { + addHandler: vi.fn((handler: () => Promise) => settingLoadedHandlers.push(handler)), + }, + }, + context, + database: { isDatabaseReady: vi.fn(() => true) }, + databaseEvents: { onDatabaseInitialised: { addHandler: vi.fn() } }, + keyValueDB: { kvDB: keyValueDB }, + localDatabase, + path: { getPath: vi.fn((entry: { path: string }) => entry.path) }, + replication: { + databaseQueueCount: { value: 0 }, + storageApplyingCount: { value: 0 }, + replicationResultCount: { value: 0 }, + onBeforeReplicate: { addHandler: vi.fn() }, + onPrepareCentralRemoteReplication: { addHandler: vi.fn() }, + onReplicationFailed: { addHandler: vi.fn() }, + parseSynchroniseResult: { + addHandler: vi.fn((handler: ParseHandler) => { + parseHandler = handler; + }), + }, + processOptionalSynchroniseResult: vi.fn(async () => false), + processSynchroniseResult, + processVirtualDocument: vi.fn(async () => false), + replicateUnattendedByEvent: vi.fn(async () => ({ status: "completed" as const })), + }, + replicator: { + createRemoteResource: vi.fn(async () => ({ + read: vi.fn(async () => new Uint8Array([1])), + dispose: vi.fn(), + })), + onBeforeReplicatorPublication: { addHandler: vi.fn() }, + onCloseActiveReplication: vi.fn(async () => true), + }, + setting: { currentSettings: vi.fn(() => settings) }, + tweakValue: {}, + vault: { + isFileSizeTooLarge: vi.fn(() => false), + isTargetFile: vi.fn(async () => true), + isValidPath: vi.fn(() => true), + }, + }; + const core = { + confirm: {}, + get localDatabase() { + return localDatabase; + }, + rebuilder: {}, + services, + }; + + useReplicationFeature(core as never); + + return { + context, + get applicationReady() { + return applicationReady; + }, + processSynchroniseResult, + get parseHandler() { + return parseHandler; + }, + settingLoadedHandlers, + setApplicationReady(value: boolean) { + applicationReady = value; + }, + settings, + }; +} + +describe("received change readiness composition", () => { + it("applies a queued received document once readiness is established and preserves explicit suspension", async () => { + eventHub.offAll(); + const harness = setup(); + try { + await harness.settingLoadedHandlers[0](); + + await harness.parseHandler!([receivedNote("ready-note")]); + await new Promise((resolve) => setTimeout(resolve, 0)); + expect(harness.processSynchroniseResult).not.toHaveBeenCalled(); + + harness.setApplicationReady(true); + harness.context.events.emitEvent(EVENT_APPLICATION_READY); + harness.context.events.emitEvent(EVENT_APPLICATION_READY); + + await vi.waitFor(() => expect(harness.processSynchroniseResult).toHaveBeenCalledTimes(1)); + + harness.settings.suspendParseReplicationResult = true; + eventHub.emitEvent(EVENT_SETTING_SAVED, harness.settings as never); + harness.setApplicationReady(false); + await harness.parseHandler!([receivedNote("suspended-note")]); + await new Promise((resolve) => setTimeout(resolve, 0)); + expect(harness.processSynchroniseResult).toHaveBeenCalledTimes(1); + + harness.setApplicationReady(true); + harness.context.events.emitEvent(EVENT_APPLICATION_READY); + harness.context.events.emitEvent(EVENT_APPLICATION_READY); + await new Promise((resolve) => setTimeout(resolve, 0)); + expect(harness.processSynchroniseResult).toHaveBeenCalledTimes(1); + + harness.settings.suspendParseReplicationResult = false; + eventHub.emitEvent(EVENT_SETTING_SAVED, harness.settings as never); + + await vi.waitFor(() => expect(harness.processSynchroniseResult).toHaveBeenCalledTimes(2)); + expect(harness.processSynchroniseResult).toHaveBeenLastCalledWith( + expect.objectContaining({ _id: "suspended-note", path: "suspended-note.md" }) + ); + } finally { + eventHub.offAll(); + } + }); +}); diff --git a/src/serviceFeatures/replication/replicationFeature.unit.spec.ts b/src/serviceFeatures/replication/replicationFeature.unit.spec.ts index fc87dc7a..eb3ff5f3 100644 --- a/src/serviceFeatures/replication/replicationFeature.unit.spec.ts +++ b/src/serviceFeatures/replication/replicationFeature.unit.spec.ts @@ -3,7 +3,9 @@ import { createServiceContext } from "@vrtmrz/livesync-commonlib/context"; import { VERSIONING_DOCID, type EntryDoc } from "@vrtmrz/livesync-commonlib/compat/common/types"; import { REMOTE_FEATURE_GENERATION } from "@vrtmrz/livesync-commonlib/replication"; import { promiseWithResolvers } from "octagonal-wheels/promises"; +import { EVENT_APPLICATION_READY } from "@/common/events"; import { useReplicationFeature } from "./index"; +import { ReplicateResultProcessor } from "./ReplicateResultProcessor"; type BooleanHandler = (showMessage: boolean) => Promise; type ParseHandler = (documents: PouchDB.Core.ExistingDocument[]) => Promise; @@ -95,6 +97,7 @@ function setup(options: SetupOptions = {}) { return { beforeReplicateHandlers, centralRemoteHandlers, + context: services.context, createRemoteResource, dispose, get parseHandler() { @@ -168,6 +171,22 @@ describe("replication serviceFeature composition", () => { expect(createRemoteResource).toHaveBeenCalledOnce(); }); + it("continues held results when Commonlib establishes application readiness", () => { + const continueHeldDocuments = vi + .spyOn(ReplicateResultProcessor.prototype, "continueHeldDocuments") + .mockImplementation(() => undefined); + try { + const { context } = setup(); + expect(continueHeldDocuments).not.toHaveBeenCalled(); + + context.events.emitEvent(EVENT_APPLICATION_READY); + + expect(continueHeldDocuments).toHaveBeenCalledOnce(); + } finally { + continueHeldDocuments.mockRestore(); + } + }); + it("requests owner retirement without awaiting the transition from result application", async () => { const retirement = promiseWithResolvers(); const onCloseActiveReplication = vi.fn(() => retirement.promise); diff --git a/test/e2e-obsidian/README.md b/test/e2e-obsidian/README.md index 37e83f17..62608bbe 100644 --- a/test/e2e-obsidian/README.md +++ b/test/e2e-obsidian/README.md @@ -221,6 +221,8 @@ This proves in real Obsidian the plug-in behaviour shared by supported platforms The workflow also opens the Customisation Sync dialogue. `--case=mtime` gives the multi-file plug-in fixture modern millisecond timestamps and requires an older remote copy to remain labelled **Older** and unselected by **Select All Shiny**. With no case argument, this check runs before the existing apply, update, and deletion workflow. +`test:e2e:obsidian:received-change-readiness` uses two sequential real Obsidian sessions and isolated CouchDB databases. The source creates ordinary notes and their Chunks; the target starts continuous replication, resets readiness through the public lifecycle service, and receives those documents while its Vault files remain absent. Marking the target ready twice must emit one readiness event and reflect the queued content. A second note remains absent after readiness while database reflecting is explicitly suspended, then appears when that setting is resumed. This focused scenario is outside `test:e2e:obsidian:local-suite`. + `test:e2e:obsidian:remote-feature-change` starts real Obsidian with continuous CouchDB replication, then changes the remote version document from generation 12 to generation 13 with an unknown feature. It waits for the control document to reach the local database and the active Replicator to retire, checks that another replication is refused, and verifies that an already accepted Vault note remains intact. After restarting the same Vault, the current remote declaration still blocks finite replication and the actual continuous connection attempt. No KV feature history is involved. It is a focused test outside `test:e2e:obsidian:local-suite`; recovery with a future compatible client remains a separate validation boundary. `test:e2e:obsidian:internal-metadata-migration` enables internal Metadata encryption through the settings UI without Rebuild. It checks unchanged plaintext and rewritten encrypted Hidden File Sync and Customisation Sync Metadata in CouchDB, stable document IDs, mismatch rejection on a second device, and file restoration after aligning settings. It then turns the preference OFF, runs Fast Fetch, and compares content loaded from both Metadata representations and their Chunks while retaining the remote feature declaration. These focused tests use the local CouchDB fixture and are outside `test:e2e:obsidian:local-suite`. diff --git a/test/e2e-obsidian/scripts/received-change-readiness.ts b/test/e2e-obsidian/scripts/received-change-readiness.ts new file mode 100644 index 00000000..690c363f --- /dev/null +++ b/test/e2e-obsidian/scripts/received-change-readiness.ts @@ -0,0 +1,420 @@ +import { readFile } from "node:fs/promises"; +import { join } from "node:path"; +import { EVENT_APPLICATION_READY } from "@vrtmrz/livesync-commonlib/compat/events/coreEvents"; +import { evalObsidianJson } from "../runner/cli.ts"; +import { + assertCouchDbReachable, + createCouchDbDatabase, + deleteCouchDbDatabase, + fetchCouchDbDocument, + loadCouchDbConfig, + makeUniqueDatabaseName, + putCouchDbDocument, + waitForCouchDbDocs, + type CouchDbConfig, + type CouchDbDocument, +} from "../runner/couchdb.ts"; +import { discoverObsidianCli, requireObsidianBinary } from "../runner/environment.ts"; +import { + assertEqual, + createE2eCouchDbPluginData, + createE2eObsidianDeviceLocalState, + prepareRemote, + pushLocalChanges, + waitForLiveSyncCoreReady, + waitForLocalDatabaseEntry, + type LocalDatabaseEntry, +} from "../runner/liveSyncWorkflow.ts"; +import { startObsidianLiveSyncSession, type ObsidianLiveSyncSession } from "../runner/session.ts"; +import { createTemporaryVault, type TemporaryVault } from "../runner/vault.ts"; + +process.env.E2E_OBSIDIAN_COUCHDB_TIMEOUT_MS ??= "20000"; + +const observerKey = "__livesyncE2eReceivedChangeReadiness"; +const readyPath = "E2E/received-change-readiness/ready.md"; +const suspendedPath = "E2E/received-change-readiness/suspended.md"; +const readyContent = [ + "# Received before readiness", + "", + "This note is replicated into the target database before the application readiness event.", + "Its content is longer than one configured chunk so the test uses the ordinary Chunk path.", + "", +].join("\n"); +const suspendedContent = [ + "# Received during explicit suspension", + "", + "This note remains queued while database reflecting is explicitly suspended.", + "Resuming result application must reflect the received metadata and its stored Chunks.", + "", +].join("\n"); + +type ReadinessObservation = { + readinessEvents: number; + content: string | null; +}; + +type CapturedNote = { + entry: LocalDatabaseEntry; + documents: CouchDbDocument[]; +}; + +async function waitForVaultContent( + vault: TemporaryVault, + path: string, + expected: string, + timeoutMs = Number(process.env.E2E_OBSIDIAN_FILE_TIMEOUT_MS ?? 10000) +): Promise { + const fullPath = join(vault.path, path); + const deadline = Date.now() + timeoutMs; + let lastContent: string | null = null; + while (Date.now() < deadline) { + try { + lastContent = await readFile(fullPath, "utf8"); + if (lastContent === expected) return; + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== "ENOENT") throw error; + } + await new Promise((resolve) => setTimeout(resolve, 100)); + } + throw new Error(`Timed out waiting for ${path} in the target Vault. Last content: ${String(lastContent)}`); +} + +async function assertVaultPathStaysMissing(vault: TemporaryVault, path: string, durationMs: number): Promise { + const fullPath = join(vault.path, path); + const deadline = Date.now() + durationMs; + while (Date.now() < deadline) { + try { + const content = await readFile(fullPath, "utf8"); + throw new Error(`The suspended or unready document was reflected early at ${path}: ${content}`); + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== "ENOENT") throw error; + } + await new Promise((resolve) => setTimeout(resolve, 100)); + } +} + +async function createAndUploadNote( + cliBinary: string, + env: NodeJS.ProcessEnv, + couchDb: CouchDbConfig, + dbName: string, + path: string, + content: string +): Promise { + const created = await evalObsidianJson<{ path: string; content: string }>( + cliBinary, + `(async()=>{const path=${JSON.stringify(path)};const folders=path.split('/');for(let i=1;i { + const ids = new Set(documents.map((document) => document._id)); + return ids.has(entry.id) && entry.children.every((childId) => ids.has(childId)); + }); + const documents = await Promise.all( + [...new Set([entry.id, ...entry.children])].map((documentId) => + fetchCouchDbDocument(couchDb, dbName, documentId) + ) + ); + return { entry, documents }; +} + +async function injectCapturedNote( + couchDb: CouchDbConfig, + dbName: string, + note: CapturedNote, + publishedChunkIds: Set +): Promise { + const chunks = note.documents.filter((document) => document._id !== note.entry.id); + const metadata = note.documents.find((document) => document._id === note.entry.id); + if (!metadata) throw new Error(`The source note metadata was not captured: ${note.entry.path}`); + for (const document of chunks) { + if (publishedChunkIds.has(document._id)) continue; + const freshDocument = { ...document }; + delete freshDocument._rev; + await putCouchDbDocument(couchDb, dbName, freshDocument); + publishedChunkIds.add(document._id); + } + const freshMetadata = { ...metadata }; + delete freshMetadata._rev; + await putCouchDbDocument(couchDb, dbName, freshMetadata); +} + +async function installObserver(cliBinary: string, env: NodeJS.ProcessEnv): Promise { + await evalObsidianJson( + cliBinary, + [ + "(()=>{", + `const key=${JSON.stringify(observerKey)};`, + `const readyEvent=${JSON.stringify(EVENT_APPLICATION_READY)};`, + "const previous=globalThis[key];", + "if(previous) previous.unsubscribeReady?.();", + "const state={readinessEvents:0,unsubscribeReady:null};", + "state.unsubscribeReady=app.plugins.plugins['obsidian-livesync'].core.services.context.events.onEvent(readyEvent,()=>state.readinessEvents++);", + "globalThis[key]=state;", + "return JSON.stringify(true);", + "})()", + ].join(""), + env + ); +} + +async function readObservation(cliBinary: string, env: NodeJS.ProcessEnv, path: string): Promise { + return await evalObsidianJson( + cliBinary, + [ + "(async()=>{", + `const key=${JSON.stringify(observerKey)};`, + `const path=${JSON.stringify(path)};`, + "const state=globalThis[key];", + "const file=app.vault.getAbstractFileByPath(path);", + "return JSON.stringify({readinessEvents:state?.readinessEvents??0,content:file?await app.vault.read(file):null});", + "})()", + ].join(""), + env + ); +} + +async function resetReadiness(cliBinary: string, env: NodeJS.ProcessEnv, suspendReflecting: boolean): Promise { + const result = await evalObsidianJson<{ ready: boolean }>( + cliBinary, + [ + "(async()=>{", + "const core=app.plugins.plugins['obsidian-livesync'].core;", + ...(suspendReflecting + ? [ + "await core.services.setting.applyPartial({suspendParseReplicationResult:true},true);", + "await core.services.control.applySettings();", + "await core.services.replication.startContinuous({trigger:'daemon',interaction:{kind:'forbidden'}});", + ] + : []), + "core.services.appLifecycle.resetIsReady();", + "return JSON.stringify({ready:core.services.appLifecycle.isReady()});", + "})()", + ].join(""), + env + ); + assertEqual(result.ready, false, "The target application remained ready after resetIsReady()."); +} + +async function waitForReceivedWhileUnready( + cliBinary: string, + env: NodeJS.ProcessEnv, + target: TemporaryVault, + path: string +): Promise { + const entry = await waitForLocalDatabaseEntry(cliBinary, env, path); + if (entry.children.length === 0) throw new Error(`The target received metadata without Chunks: ${path}`); + const readiness = await evalObsidianJson<{ ready: boolean }>( + cliBinary, + "(()=>JSON.stringify({ready:app.plugins.plugins['obsidian-livesync'].core.services.appLifecycle.isReady()}))()", + env + ); + assertEqual(readiness.ready, false, `The target became ready before applying ${path}.`); + await assertVaultPathStaysMissing(target, path, 750); +} + +async function startContinuousReplication(cliBinary: string, env: NodeJS.ProcessEnv): Promise { + const result = await evalObsidianJson<{ status: string }>( + cliBinary, + [ + "(async()=>{", + "const core=app.plugins.plugins['obsidian-livesync'].core;", + "await core.services.setting.applyExternalSettings({liveSync:true},true);", + "await core.services.control.applySettings();", + "const result=await core.services.replication.startContinuous({trigger:'daemon',interaction:{kind:'forbidden'}});", + "return JSON.stringify(result);", + "})()", + ].join(""), + env + ); + assertEqual(result.status, "completed", `Continuous replication did not start: ${JSON.stringify(result)}.`); + + const deadline = Date.now() + Number(process.env.E2E_OBSIDIAN_REMOTE_ACTIVITY_TIMEOUT_MS ?? 30000); + let active = false; + while (!active && Date.now() < deadline) { + active = await evalObsidianJson( + cliBinary, + "(()=>JSON.stringify(!!app.plugins.plugins['obsidian-livesync'].core.services.replicator.getActiveReplicator()))()", + env + ); + if (!active) await new Promise((resolve) => setTimeout(resolve, 250)); + } + if (!active) throw new Error("Timed out waiting for the target Replicator to become active."); +} + +async function markReadyTwice(cliBinary: string, env: NodeJS.ProcessEnv): Promise { + const result = await evalObsidianJson<{ ready: boolean }>( + cliBinary, + [ + "(()=>{", + "const lifecycle=app.plugins.plugins['obsidian-livesync'].core.services.appLifecycle;", + "lifecycle.markIsReady();", + "lifecycle.markIsReady();", + "return JSON.stringify({ready:lifecycle.isReady()});", + "})()", + ].join(""), + env + ); + assertEqual(result.ready, true, "Commonlib did not establish application readiness."); +} + +async function removeObserver(cliBinary: string, env: NodeJS.ProcessEnv): Promise { + await evalObsidianJson( + cliBinary, + [ + "(()=>{", + `const key=${JSON.stringify(observerKey)};`, + "const state=globalThis[key];", + "if(state){state.unsubscribeReady?.();delete globalThis[key];}", + "return JSON.stringify(true);", + "})()", + ].join(""), + env + ); +} + +async function assertApplied( + cliBinary: string, + env: NodeJS.ProcessEnv, + vault: TemporaryVault, + path: string, + expectedContent: string, + expectedReadyEvents: number +): Promise { + const beforeReflection = await readObservation(cliBinary, env, path); + assertEqual( + beforeReflection.readinessEvents, + expectedReadyEvents, + "The expected Commonlib readiness transition was not observed before reflection." + ); + await waitForVaultContent(vault, path, expectedContent); + await new Promise((resolve) => setTimeout(resolve, 500)); + const observation = await readObservation(cliBinary, env, path); + assertEqual(observation.content, expectedContent, `Obsidian Vault read-back was incorrect for ${path}.`); + assertEqual( + observation.readinessEvents, + expectedReadyEvents, + "markIsReady emitted an unexpected number of readiness transitions." + ); +} + +async function main(): Promise { + const binary = requireObsidianBinary(); + const cli = discoverObsidianCli(); + if (!cli.binary) throw new Error(`Could not find obsidian-cli. Checked paths: ${cli.checked.join(", ")}`); + + const couchDb = await loadCouchDbConfig(); + const sourceDbName = makeUniqueDatabaseName(couchDb.dbPrefix, "received-change-readiness-source"); + const targetDbName = makeUniqueDatabaseName(couchDb.dbPrefix, "received-change-readiness-target"); + const couchDbSettings = { + uri: couchDb.uri, + username: couchDb.username, + password: couchDb.password, + }; + const sourceCouchDbSettings = { ...couchDbSettings, dbName: sourceDbName }; + const targetCouchDbSettings = { ...couchDbSettings, dbName: targetDbName }; + const sourceVault = await createTemporaryVault("obsidian-livesync-readiness-source-"); + const targetVault = await createTemporaryVault("obsidian-livesync-readiness-target-"); + let source: ObsidianLiveSyncSession | undefined; + let target: ObsidianLiveSyncSession | undefined; + const publishedChunkIds = new Set(); + + try { + await assertCouchDbReachable(couchDb); + await createCouchDbDatabase(couchDb, sourceDbName); + await createCouchDbDatabase(couchDb, targetDbName); + source = await startObsidianLiveSyncSession({ + binary, + cliBinary: cli.binary, + vault: sourceVault, + startupGraceMs: Number(process.env.E2E_OBSIDIAN_STARTUP_GRACE_MS ?? 1000), + pluginData: createE2eCouchDbPluginData(sourceCouchDbSettings), + localStorageEntries: createE2eObsidianDeviceLocalState(sourceVault.name), + }); + await waitForLiveSyncCoreReady(cli.binary, source.cliEnv); + await prepareRemote(cli.binary, source.cliEnv); + + const readyNote = await createAndUploadNote( + cli.binary, + source.cliEnv, + couchDb, + sourceDbName, + readyPath, + readyContent + ); + const suspendedNote = await createAndUploadNote( + cli.binary, + source.cliEnv, + couchDb, + sourceDbName, + suspendedPath, + suspendedContent + ); + await source.app.stop(); + source = undefined; + + target = await startObsidianLiveSyncSession({ + binary, + cliBinary: cli.binary, + vault: targetVault, + startupGraceMs: Number(process.env.E2E_OBSIDIAN_STARTUP_GRACE_MS ?? 1000), + pluginData: createE2eCouchDbPluginData(targetCouchDbSettings), + localStorageEntries: createE2eObsidianDeviceLocalState(targetVault.name), + }); + await waitForLiveSyncCoreReady(cli.binary, target.cliEnv); + await prepareRemote(cli.binary, target.cliEnv); + await installObserver(cli.binary, target.cliEnv); + await startContinuousReplication(cli.binary, target.cliEnv); + + await resetReadiness(cli.binary, target.cliEnv, false); + await injectCapturedNote(couchDb, targetDbName, readyNote, publishedChunkIds); + await waitForReceivedWhileUnready(cli.binary, target.cliEnv, targetVault, readyPath); + await markReadyTwice(cli.binary, target.cliEnv); + await assertApplied(cli.binary, target.cliEnv, targetVault, readyPath, readyContent, 1); + + await resetReadiness(cli.binary, target.cliEnv, true); + await injectCapturedNote(couchDb, targetDbName, suspendedNote, publishedChunkIds); + await waitForReceivedWhileUnready(cli.binary, target.cliEnv, targetVault, suspendedPath); + await markReadyTwice(cli.binary, target.cliEnv); + await assertVaultPathStaysMissing(targetVault, suspendedPath, 1000); + const suspendedObservation = await readObservation(cli.binary, target.cliEnv, suspendedPath); + assertEqual(suspendedObservation.readinessEvents, 2, "The resumed-ready transition was not observed once."); + assertEqual(suspendedObservation.content, null, "Readiness bypassed explicit database-reflection suspension."); + + await evalObsidianJson( + cli.binary, + "(async()=>{const setting=app.plugins.plugins['obsidian-livesync'].core.services.setting;await setting.applyPartial({suspendParseReplicationResult:false},true);return JSON.stringify(true);})()", + target.cliEnv + ); + await assertApplied(cli.binary, target.cliEnv, targetVault, suspendedPath, suspendedContent, 2); + + console.log( + "Received-change readiness: queued changes resumed once per readiness transition; explicit suspension held until settings resumed." + ); + } finally { + if (target) { + await removeObserver(cli.binary, target.cliEnv).catch(() => undefined); + await target.app.stop(); + } + if (source) await source.app.stop(); + await sourceVault.dispose(); + await targetVault.dispose(); + if (process.env.E2E_OBSIDIAN_KEEP_COUCHDB !== "true") { + await Promise.all( + [sourceDbName, targetDbName].map((dbName) => deleteCouchDbDatabase(couchDb, dbName)) + ).catch((error: unknown) => { + console.warn(error instanceof Error ? error.message : error); + }); + } + } +} + +main().catch((error: unknown) => { + console.error(error instanceof Error ? error.stack : error); + process.exit(1); +}); diff --git a/test/e2e-obsidian/scripts/run-focused.ts b/test/e2e-obsidian/scripts/run-focused.ts index ddf2e049..e3d2691e 100644 --- a/test/e2e-obsidian/scripts/run-focused.ts +++ b/test/e2e-obsidian/scripts/run-focused.ts @@ -33,6 +33,7 @@ const focusedScenarios = new Set([ "security-seed-reconnect", "hidden-file-snippet-sync", "customisation-sync", + "received-change-readiness", "setting-markdown-export", "upgrade-from-stable", ]); diff --git a/updates.md b/updates.md index f794934f..e036efa9 100644 --- a/updates.md +++ b/updates.md @@ -37,6 +37,12 @@ Earlier releases remain available in the 1.0 release history, the 1.0 preview hi - Retries start after two seconds and continue with increasing delays while finite replication is active. When it ends, LiveSync checks locally and makes a final lookup if needed, without waiting out the remaining retry delay. - We can now distinguish initial on-demand Chunk requests (`🛄`) from retries (`🔁`) in the status bar. These replace `🧩`; each pending Chunk appears in one category, including while a retry is waiting. +### Synchronisation and storage + +#### Fixed + +- Received changes held during start-up or a fetch are applied when LiveSync becomes ready, without waiting for another change or a settings save. **Suspend database reflecting** continues to hold changes (#1200). + ## 1.0.32 27th September, 2026