Resume held replication results on application readiness

Trigger queued result processing on Commonlib readiness while preserving
explicit suspension. Cover both transitions with focused tests.
This commit is contained in:
Kim Jansheden
2026-09-29 19:39:20 +02:00
parent b3208f2aaf
commit de2ebaeea5
5 changed files with 66 additions and 1 deletions
@@ -108,8 +108,15 @@ export class ReplicateResultProcessor {
} }
public resume() { public resume() {
this._suspended = false; 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(); this.updateProcessingActivity();
fireAndForget(() => this.runProcessQueue()); this.triggerProcessQueue();
} }
// Whether the processing is suspended // Whether the processing is suspended
@@ -207,6 +207,35 @@ describe("ReplicateResultProcessor", () => {
expect(isReady).toHaveBeenCalledOnce(); 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", () => { it("applies results in remediation mode, which never reports readiness", () => {
const { processor } = setup({ const { processor } = setup({
applicationReady: false, applicationReady: false,
+4
View File
@@ -4,6 +4,7 @@ import { UnresolvedErrorManager } from "@vrtmrz/livesync-commonlib/compat/servic
import type { ServiceContext } from "@vrtmrz/livesync-commonlib/context"; import type { ServiceContext } from "@vrtmrz/livesync-commonlib/context";
import { fireAndForget } from "octagonal-wheels/promises"; import { fireAndForget } from "octagonal-wheels/promises";
import type { IMinimumLiveSyncCommands, LiveSyncBaseCore } from "@/LiveSyncBaseCore"; import type { IMinimumLiveSyncCommands, LiveSyncBaseCore } from "@/LiveSyncBaseCore";
import { EVENT_APPLICATION_READY } from "@/common/events";
import { createAutomaticReplicationTriggers } from "./automaticTriggers"; import { createAutomaticReplicationTriggers } from "./automaticTriggers";
import { createCentralCompatibilityRecovery } from "./centralCompatibilityRecovery"; import { createCentralCompatibilityRecovery } from "./centralCompatibilityRecovery";
import { createOnlineReplicationPreflight, createSecuritySeedPreflight } from "./preflight"; import { createOnlineReplicationPreflight, createSecuritySeedPreflight } from "./preflight";
@@ -102,6 +103,9 @@ export function useReplicationFeature<TContext extends ServiceContext, TCommands
fireAndForget(() => resultProcessor.restoreFromSnapshotOnce()); fireAndForget(() => resultProcessor.restoreFromSnapshotOnce());
return Promise.resolve(true); 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.appLifecycle.onSettingLoaded.addHandler(initialiseAutomaticReplicationTriggers);
services.replication.parseSynchroniseResult.addHandler((documents) => { services.replication.parseSynchroniseResult.addHandler((documents) => {
resultProcessor.enqueueAll(documents); resultProcessor.enqueueAll(documents);
@@ -3,7 +3,9 @@ import { createServiceContext } from "@vrtmrz/livesync-commonlib/context";
import { VERSIONING_DOCID, type EntryDoc } from "@vrtmrz/livesync-commonlib/compat/common/types"; import { VERSIONING_DOCID, type EntryDoc } from "@vrtmrz/livesync-commonlib/compat/common/types";
import { REMOTE_FEATURE_GENERATION } from "@vrtmrz/livesync-commonlib/replication"; import { REMOTE_FEATURE_GENERATION } from "@vrtmrz/livesync-commonlib/replication";
import { promiseWithResolvers } from "octagonal-wheels/promises"; import { promiseWithResolvers } from "octagonal-wheels/promises";
import { EVENT_APPLICATION_READY } from "@/common/events";
import { useReplicationFeature } from "./index"; import { useReplicationFeature } from "./index";
import { ReplicateResultProcessor } from "./ReplicateResultProcessor";
type BooleanHandler = (showMessage: boolean) => Promise<boolean>; type BooleanHandler = (showMessage: boolean) => Promise<boolean>;
type ParseHandler = (documents: PouchDB.Core.ExistingDocument<EntryDoc>[]) => Promise<boolean>; type ParseHandler = (documents: PouchDB.Core.ExistingDocument<EntryDoc>[]) => Promise<boolean>;
@@ -95,6 +97,7 @@ function setup(options: SetupOptions = {}) {
return { return {
beforeReplicateHandlers, beforeReplicateHandlers,
centralRemoteHandlers, centralRemoteHandlers,
context: services.context,
createRemoteResource, createRemoteResource,
dispose, dispose,
get parseHandler() { get parseHandler() {
@@ -168,6 +171,22 @@ describe("replication serviceFeature composition", () => {
expect(createRemoteResource).toHaveBeenCalledOnce(); 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 () => { it("requests owner retirement without awaiting the transition from result application", async () => {
const retirement = promiseWithResolvers<boolean>(); const retirement = promiseWithResolvers<boolean>();
const onCloseActiveReplication = vi.fn(() => retirement.promise); const onCloseActiveReplication = vi.fn(() => retirement.promise);
+6
View File
@@ -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) - 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. - 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 ## 1.0.32
27th September, 2026 27th September, 2026