From de2ebaeea5d29c53c05bcb32db0ec70f93b6428d Mon Sep 17 00:00:00 2001 From: Kim Jansheden Date: Tue, 29 Sep 2026 19:39:20 +0200 Subject: [PATCH] Resume held replication results on application readiness Trigger queued result processing on Commonlib readiness while preserving explicit suspension. Cover both transitions with focused tests. --- .../replication/ReplicateResultProcessor.ts | 9 +++++- .../ReplicateResultProcessor.unit.spec.ts | 29 +++++++++++++++++++ src/serviceFeatures/replication/index.ts | 4 +++ .../replicationFeature.unit.spec.ts | 19 ++++++++++++ updates.md | 6 ++++ 5 files changed, 66 insertions(+), 1 deletion(-) 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/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/updates.md b/updates.md index f11a6b45..e9becde8 100644 --- a/updates.md +++ b/updates.md @@ -34,6 +34,12 @@ Earlier releases remain available in the 1.0 release history, the 1.0 preview hi - We can now keep using an E2EE passphrase beginning with `%` after restarting Obsidian. (#1221) - LiveSync encrypts it before saving the settings. If an earlier version saved it in plain text, re-enter the passphrase used to encrypt the existing data after updating. Treat that passphrase as exposed if the affected `data.json` was shared. +### 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