Simplify remote feature checks and clarify manual rebuild guidance

This commit is contained in:
vorotamoroz
2026-09-27 11:08:53 +00:00
parent 9a047d44f3
commit 2c6cd56618
15 changed files with 251 additions and 684 deletions
@@ -75,70 +75,44 @@ for compatible clients in the explanation. Do not set `requireRebuild` or
operations and restart. `recommendRebuild` currently exists only as an unused operations and restart. `recommendRebuild` currently exists only as an unused
rule field, so setting it alone does not display an explanation. rule field, so setting it alone does not display an explanation.
## Receiving a changed version document ## Admission and received version documents
The existing path is: The remote version document is the source of feature requirements. Commonlib
checks it before replication, Fast Fetch, and direct access. The milestone keeps
the existing Tweak comparison and Rebuild lock. An accepted writer declares the
feature before using it, including the writer admitted to a locked rebuilt
remote; an unaccepted device remains blocked by that lock.
1. Commonlib receives a CouchDB replication change and calls Retain the existing received-version path through `parseSynchroniseResult`,
`parseSynchroniseResult` with the received documents. `enqueueAll`, and `processIfNonDocumentChange`. Replace its numeric comparison
2. The replication service feature passes them to with the shared assessment so unknown names at the same generation are also
`ReplicateResultProcessor.enqueueAll`. reported. Known features, reordered lists, and ordinary revision updates do not
3. `processIfNonDocumentChange` recognises `type: versioninfo` and requests active retire the connection. Unsupported or malformed control documents request
Replicator retirement when `version > VER`. retirement through the existing Replicator owner and display the reason.
4. The owner closes admission, requests transfer cancellation, drains its work, The callback must not await retirement of the operation which delivered it.
and closes the instance. The result callback does not wait for that transition.
Retain this observation path, but use Commonlib's complete assessment of the This is an admission check and a best-effort stop for exceptional changes during
identified control document. A changed `used_features` list must be inspected an active connection. It does not fence every queued file application, roll back
even when the numeric version is unchanged. A mere revision change, list reorder, accepted writes, or guarantee an atomic change across live devices. Feature
or duplicate known identifier does not require an interruption. changes are an infrequent administrative operation: update all devices first,
then enable the preference and use the recommended manual Rebuild. Rebuild
Inspect the entire batch's control information before passing any file entries locks the remote using the existing workflow; changing this preference alone
to normal or optional processing. Recognise the fixed control-document ID and does not lock it. The action to proceed without rebuilding explicitly reminds
validate its type and contents. A deleted or malformed control document is a the person to update every other device, including currently connected devices.
rejection, not a successful empty update.
When all requirements remain supported, refresh the assessment and affected
shared-setting checks; a newly added feature retires the current writer so its
next admission rechecks the shared Tweak policy. When a requirement is
unknown or incompatible, synchronously record the block for the affected
database, stop admitting new reflection and database operations, and request
retirement through the existing owner. Notify with the unknown identifiers as
text, using a generic message when no descriptive label exists.
File application and remote transfer have separate lifetimes. Requesting owner
retirement alone is not the application block. Keep the block separate from
temporary lifecycle suspension so an ordinary resume event cannot clear it.
Queued or waiting work checks it before starting another write; notifications
from an old physical database must not affect its replacement.
Do not await `onCloseActiveReplication` inside the callback which delivered the
change. That callback can belong to work which retirement must drain. Establish
the block immediately, request retirement without awaiting it there, and let the
owner perform cancellation and close in its existing order.
## Persistence and recovery boundaries ## Persistence and recovery boundaries
A CouchDB replication notification can arrive after the documents have entered Do not retain a second feature list, highest generation, or rejection flag in
the local DB. Already-started network and filesystem operations may settle. KV storage. Do not add compatibility checks to pending-work snapshot recovery
This feature does not promise rollback or atomic revocation of those operations. or make that recovery a new prerequisite for application readiness. Preserve
the existing queue and startup behaviour. A later attempt checks the current
remote declaration, including after restart. Declared features remain on the
remote when the write preference is disabled because older data can still use
them; manually shortening that declaration is not a supported migration.
Preserve pending work or durable reconciliation information when stopping. A After updating clients, use normal reconnection and the existing Hatch
checkpoint may already include the documents which have not reached the Vault. inspection or Fetch workflow if reconciliation is needed. This feature does not
Do not drop those documents or depend on an ordinary reconnect to send them repair unrelated KV inconsistencies or the existing readiness queue behaviour.
again. A compatible client must reassess and reprocess or explicitly reacquire
them before lifting the applicable block.
Persist a blocked pending-work snapshot before the received-change callback
settles, while requesting owner retirement separately to avoid a circular wait.
Restore compatibility checks before replication result application and the next
ordinary synchronisation. Retain observed feature requirements with the pending
work snapshot so that a shortened version-document list does not release a
blocked local database on restart. A dismissed Notice or a changed connection
does not establish that the affected local data has become interpretable.
An older local generation without feature declarations remains readable during
normal application and cleaned-remote recovery; remote migration remains the
responsibility of the replication admission check.
Garbage Collection V3 is a beta manual operation which begins with an ordinary Garbage Collection V3 is a beta manual operation which begins with an ordinary
bidirectional synchronisation. That admission checks the remote feature bidirectional synchronisation. That admission checks the remote feature
@@ -156,40 +130,18 @@ Keep focused tests for settings defaults and imports, the Doctor condition
matrix, acceptance and dismissal, connection replacement, and absence of an matrix, acceptance and dismissal, connection replacement, and absence of an
automatic Rebuild, Fetch, or restart for this rule. automatic Rebuild, Fetch, or restart for this rule.
Add deterministic runtime tests for feature-only changes, version documents Keep unit tests for known and unknown feature notifications, generic identifier
first and last in a batch, unknown-name presentation, duplicate notifications, presentation, retirement without a circular wait, and the unchanged snapshot
queued and waiting reflection, stale database callbacks, restart, checkpointed behaviour after KV failure or obsolete snapshot fields. The previous batch
but unapplied documents, and cancellation without a circular wait. fences, physical-database tracking, and persistent rejection tests are outside
this design; they must not imply an atomic live migration guarantee.
Before this implementation, the host processor was checked with a focused unit probe: a numeric Use real Obsidian Hidden File Sync and Customisation Sync scenarios to inspect
incompatibility requests retirement, but the processor still applies a note in raw CouchDB Metadata and restore content in another Vault. Check the admitted
the same batch when its host remains ready. A same-version document with an writer on a locked remote, unknown-feature rejection before and during
unknown feature does not request retirement. These observations motivated the replication, and remote-based rejection after restart. Retain the encrypted
new checks. CLI-to-Obsidian interoperability scenario. A future client upgrade that adds
support for an unknown feature is a separate validation boundary.
Extend real Obsidian Hidden File Sync and Customisation Sync scenarios and the
CLI interoperability checks. Inspect raw CouchDB documents as well as restored
files. Exercise an active connection when another client changes the feature
requirements, and verify that previously accepted data survives the stop.
Pending-work restoration and recovery with a compatible client are separate
boundaries.
The local packed Commonlib 0.1.30 candidate passed Commonlib unit and boundary
tests and the LiveSync build, type checks, and unit tests. Real Obsidian 1.12.7
passed the Hidden File Sync, Customisation Sync, and encrypted CLI-to-Obsidian
scenarios. The dedicated active-connection scenario changed a generation 12
remote to generation 13 with an unknown feature whilst continuous replication
was running. The local control document arrived, the active Replicator retired,
a subsequent replication was refused, and an earlier accepted note stayed in
the Vault. The same real Obsidian scenario now shortens both control-document
feature lists and restarts the Vault. The saved observation remains associated
with the same physical database, and both OneShot and Continuous replication
are refused. Focused host tests cover both batch orders, pending-work snapshots,
stale callbacks, and the older local-generation case. Recovery after upgrading
to a future client that understands the unknown feature has not been exercised
in real Obsidian. The existing [readiness queue issue](https://github.com/vrtmrz/obsidian-livesync/issues/1200)
still affects when restored pending documents resume after the application
becomes ready; snapshot preservation alone does not resolve that issue.
Keep the primary-language settings and troubleshooting guides, the Keep the primary-language settings and troubleshooting guides, the
database-compatibility ADR, and Unreleased notes aligned with this behaviour. database-compatibility ADR, and Unreleased notes aligned with this behaviour.
+1 -1
View File
@@ -251,7 +251,7 @@ Setting key: encryptInternalMetadata
For CouchDB, this encrypts paths, times, sizes, and Chunk references in the Metadata used by Hidden File Sync and Customisation Sync. It requires E2EE V2 and **Property Encryption**. New Vaults enable the preference by default, but it has no effect until those prerequisites are enabled. Existing Vaults and older Setup URIs and QR codes keep it disabled unless you enable it. For CouchDB, this encrypts paths, times, sizes, and Chunk references in the Metadata used by Hidden File Sync and Customisation Sync. It requires E2EE V2 and **Property Encryption**. New Vaults enable the preference by default, but it has no effect until those prerequisites are enabled. Existing Vaults and older Setup URIs and QR codes keep it disabled unless you enable it.
Enabling the preference protects future Metadata writes. Existing Metadata and earlier revisions can remain readable in the remote database. If you want to protect existing Metadata too, prepare the authoritative data, update every device to a compatible version, and manually Rebuild the remote database. LiveSync does not gather data or start a Rebuild when you change this preference. Plaintext and encrypted Metadata can coexist during the transition. Document IDs, revision information, document counts, and ciphertext lengths remain visible. Enabling the preference protects future Metadata writes. Existing Metadata and earlier revisions can remain readable in the remote database. If you want to protect existing Metadata too, prepare the authoritative data, update every device to a compatible version, and manually Rebuild the remote database. LiveSync does not gather data or start a Rebuild when you change this preference. The action to enable it without rebuilding explicitly reminds you to update every other synchronising device first, including devices currently running LiveSync. Plaintext and encrypted Metadata can coexist during the transition. Document IDs, revision information, document counts, and ciphertext lengths remain visible.
#### Encryption Algorithm #### Encryption Algorithm
+1 -1
View File
@@ -113,7 +113,7 @@ Historic defect notices and renamed controls are retained in the [0.25 release h
## The remote database uses an unknown feature ## The remote database uses an unknown feature
When a notice identifies an unknown feature, update this device and every other client of the same CouchDB database, including the CLI. The notice includes the feature identifier even if this version has no descriptive name for it. Synchronisation and pending file reflection pause because an older client may not interpret the Metadata and its Chunk references correctly. The cleaned-remote recovery path also checks compatibility before counting Chunk references. Do not remove the feature name from the remote version document to bypass the check. After updating, reconnect and review any pending file changes before running Garbage Collection. When a notice identifies an unknown feature, update this device and every other client of the same CouchDB database, including the CLI. The notice includes the feature identifier even if this version has no descriptive name for it. New synchronisation is refused, and receiving an unsupported requirement stops active replication, because an older client may not interpret the Metadata and its Chunk references correctly. Already queued file changes are not rolled back. The cleaned-remote recovery path also checks compatibility before counting Chunk references. Do not remove the feature name from the remote version document to bypass the check. After updating, reconnect and review any pending file changes before running Garbage Collection.
## Setup and settings questions ## Setup and settings questions
@@ -5,7 +5,6 @@ import {
DEFAULT_SETTINGS, DEFAULT_SETTINGS,
LOG_LEVEL_NOTICE, LOG_LEVEL_NOTICE,
type ObsidianLiveSyncSettings, type ObsidianLiveSyncSettings,
type EncryptionSettings,
LOG_LEVEL_VERBOSE, LOG_LEVEL_VERBOSE,
} from "@vrtmrz/livesync-commonlib/compat/common/types"; } from "@vrtmrz/livesync-commonlib/compat/common/types";
import { Menu, type ButtonComponent } from "@/deps.ts"; import { Menu, type ButtonComponent } from "@/deps.ts";
@@ -32,13 +31,11 @@ import {
import { ConnectionStringParser } from "@vrtmrz/livesync-commonlib/compat/common/ConnectionString"; import { ConnectionStringParser } from "@vrtmrz/livesync-commonlib/compat/common/ConnectionString";
import type { RemoteConfigurationResult } from "@vrtmrz/livesync-commonlib/compat/common/ConnectionString"; import type { RemoteConfigurationResult } from "@vrtmrz/livesync-commonlib/compat/common/ConnectionString";
import SetupRemote from "@/modules/features/SetupWizard/dialogs/SetupRemote.svelte"; import SetupRemote from "@/modules/features/SetupWizard/dialogs/SetupRemote.svelte";
import SetupRemoteE2EE from "@/modules/features/SetupWizard/dialogs/SetupRemoteE2EE.svelte";
import SetupRemoteCouchDB from "@/modules/features/SetupWizard/dialogs/SetupRemoteCouchDB.svelte"; import SetupRemoteCouchDB from "@/modules/features/SetupWizard/dialogs/SetupRemoteCouchDB.svelte";
import SetupRemoteBucket from "@/modules/features/SetupWizard/dialogs/SetupRemoteBucket.svelte"; import SetupRemoteBucket from "@/modules/features/SetupWizard/dialogs/SetupRemoteBucket.svelte";
import type { import type {
SetupRemoteCouchDBInitialData, SetupRemoteCouchDBInitialData,
SetupRemoteCouchDBResultType, SetupRemoteCouchDBResultType,
SetupRemoteE2EEResultType,
} from "@/modules/features/SetupWizard/dialogs/setupDialogTypes.ts"; } from "@/modules/features/SetupWizard/dialogs/setupDialogTypes.ts";
import { syncActivatedRemoteSettings } from "./remoteConfigBuffer.ts"; import { syncActivatedRemoteSettings } from "./remoteConfigBuffer.ts";
@@ -119,34 +116,15 @@ export function paneRemoteConfig(
.onClick(async () => { .onClick(async () => {
const setupManager = this.core.getModule(SetupManager); const setupManager = this.core.getModule(SetupManager);
const originalSettings = getSettingsFromEditingSettings(this.editingSettings); const originalSettings = getSettingsFromEditingSettings(this.editingSettings);
const e2eeConf = await setupManager.dialogManager.openWithExplicitCancel< const applied = await setupManager.onlyE2EEConfiguration(UserMode.Update, originalSettings);
SetupRemoteE2EEResultType, if (applied) {
EncryptionSettings this.editingSettings.encryptInternalMetadata =
>(SetupRemoteE2EE, originalSettings); this.core.settings.encryptInternalMetadata;
if (e2eeConf === "cancelled") {
return;
}
const onlyInternalMetadataPreferenceChanged =
originalSettings.encryptInternalMetadata !== e2eeConf.encryptInternalMetadata &&
originalSettings.encrypt === e2eeConf.encrypt &&
originalSettings.passphrase === e2eeConf.passphrase &&
originalSettings.E2EEAlgorithm === e2eeConf.E2EEAlgorithm &&
originalSettings.usePathObfuscation === e2eeConf.usePathObfuscation;
if (onlyInternalMetadataPreferenceChanged) {
await this.services.setting.applyPartial(
{ encryptInternalMetadata: e2eeConf.encryptInternalMetadata },
true
);
this.editingSettings.encryptInternalMetadata = e2eeConf.encryptInternalMetadata;
if (this.initialSettings) { if (this.initialSettings) {
this.initialSettings.encryptInternalMetadata = e2eeConf.encryptInternalMetadata; this.initialSettings.encryptInternalMetadata =
this.core.settings.encryptInternalMetadata;
} }
this.requestUpdate(); this.requestUpdate();
} else {
await setupManager.onConfirmApplySettingsFromWizard(
{ ...originalSettings, ...e2eeConf },
UserMode.Update
);
} }
updateE2EESummary(); updateE2EESummary();
}) })
@@ -162,24 +162,15 @@ describe("paneRemoteConfig", () => {
encryptInternalMetadata: false, encryptInternalMetadata: false,
remoteConfigurations: {}, remoteConfigurations: {},
}; };
const applyPartial = vi.fn(async () => {});
const onConfirmApplySettingsFromWizard = vi.fn(async () => {});
const setupManager = { const setupManager = {
dialogManager: { onlyE2EEConfiguration: vi.fn(async () => {
openWithExplicitCancel: vi.fn(async () => ({ host.core.settings.encryptInternalMetadata = true;
encrypt: true, return true;
passphrase: "passphrase", }),
E2EEAlgorithm: "v2",
usePathObfuscation: true,
encryptInternalMetadata: true,
})),
},
onConfirmApplySettingsFromWizard,
}; };
const host = { const host = {
editingSettings: { ...originalSettings }, editingSettings: { ...originalSettings },
initialSettings: { ...originalSettings }, initialSettings: { ...originalSettings },
services: { setting: { applyPartial } },
core: { core: {
settings: { ...originalSettings }, settings: { ...originalSettings },
getModule: vi.fn(() => setupManager), getModule: vi.fn(() => setupManager),
@@ -198,8 +189,7 @@ describe("paneRemoteConfig", () => {
paneRemoteConfig.call(host as never, {} as HTMLElement, { addPanel } as never); paneRemoteConfig.call(host as never, {} as HTMLElement, { addPanel } as never);
await runtime.clickHandlers[0](); await runtime.clickHandlers[0]();
expect(applyPartial).toHaveBeenCalledWith({ encryptInternalMetadata: true }, true); expect(setupManager.onlyE2EEConfiguration).toHaveBeenCalledOnce();
expect(onConfirmApplySettingsFromWizard).not.toHaveBeenCalled();
expect(host.editingSettings.encryptInternalMetadata).toBe(true); expect(host.editingSettings.encryptInternalMetadata).toBe(true);
expect(host.initialSettings.encryptInternalMetadata).toBe(true); expect(host.initialSettings.encryptInternalMetadata).toBe(true);
expect(host.requestUpdate).toHaveBeenCalledOnce(); expect(host.requestUpdate).toHaveBeenCalledOnce();
+25
View File
@@ -336,6 +336,31 @@ export class SetupManager extends AbstractModule {
this._log("E2EE configuration cancelled.", LOG_LEVEL_NOTICE); this._log("E2EE configuration cancelled.", LOG_LEVEL_NOTICE);
return false; return false;
} }
const onlyInternalMetadataPreferenceChanged =
currentSetting.encryptInternalMetadata !== e2eeConf.encryptInternalMetadata &&
currentSetting.encrypt === e2eeConf.encrypt &&
currentSetting.passphrase === e2eeConf.passphrase &&
currentSetting.E2EEAlgorithm === e2eeConf.E2EEAlgorithm &&
currentSetting.usePathObfuscation === e2eeConf.usePathObfuscation;
if (userMode === UserMode.Update && onlyInternalMetadataPreferenceChanged) {
if (e2eeConf.encryptInternalMetadata && currentSetting.remoteType === REMOTE_COUCHDB) {
const proceed = "Enable without rebuilding — update every other device first";
const choice = await this.core.confirm.askSelectStringDialogue(
"A manual remote Rebuild is strongly recommended to protect existing Metadata. " +
"Before continuing without rebuilding, update every other synchronising device to a version " +
"which supports this option, including devices currently running LiveSync. " +
"Existing Metadata remains unchanged until it is rewritten or rebuilt.",
[proceed, "Cancel"],
{ title: "Encrypt internal file Metadata", defaultAction: "Cancel" }
);
if (choice !== proceed) return false;
}
await this.services.setting.applyPartial(
{ encryptInternalMetadata: e2eeConf.encryptInternalMetadata },
true
);
return true;
}
const newSetting = { const newSetting = {
...currentSetting, ...currentSetting,
...e2eeConf, ...e2eeConf,
@@ -659,3 +659,23 @@ describe("SetupManager", () => {
expect(setting.currentSettings().P2P_ActiveRemoteConfigurationId).toBe("existing"); expect(setting.currentSettings().P2P_ActiveRemoteConfigurationId).toBe("existing");
}); });
}); });
describe("internal Metadata configuration", () => {
it.each([true, false])(
"applies the preference only after accepting the no-Rebuild warning (%s)",
async (accept) => {
const { manager, setting, dialogManager, core } = createSetupManager();
const current = { ...setting.settings, encryptInternalMetadata: false, remoteType: REMOTE_COUCHDB };
dialogManager.openWithExplicitCancel.mockResolvedValue({ ...current, encryptInternalMetadata: true });
const ask = vi.fn(async (_message: string, choices: string[]) => (accept ? choices[0] : "Cancel"));
core.confirm = { askSelectStringDialogue: ask };
const apply = vi.spyOn(setting, "applyPartial").mockResolvedValue(undefined);
await expect(manager.onlyE2EEConfiguration(UserMode.Update, current)).resolves.toBe(accept);
expect(ask.mock.calls[0][1][0]).toContain("update every other device first");
expect(ask.mock.calls[0][0]).toContain("currently running LiveSync");
expect(apply).toHaveBeenCalledTimes(accept ? 1 : 0);
expect(core.rebuilder.scheduleRebuild).not.toHaveBeenCalled();
expect(core.rebuilder.scheduleFetch).not.toHaveBeenCalled();
}
);
});
@@ -105,7 +105,8 @@
(Obfuscate Properties). The remote type is selected later in this setup wizard. (Obfuscate Properties). The remote type is selected later in this setup wizard.
<br /> <br />
It protects Metadata written after the option is enabled; existing Metadata is not rewritten. A manual remote It protects Metadata written after the option is enabled; existing Metadata is not rewritten. A manual remote
Rebuild is strongly recommended to protect existing Metadata. Rebuild is strongly recommended to protect existing Metadata. Update every other synchronising device to a compatible
version before enabling this option, including devices currently running LiveSync.
</InfoNote> </InfoNote>
<ExtraItems title={translateMessage("Advanced")}> <ExtraItems title={translateMessage("Advanced")}>
@@ -1,7 +1,7 @@
import { assessRemoteFeatureDocument, describeRemoteFeatureRejection } from "@vrtmrz/livesync-commonlib/replication";
import { import {
SYNCINFO_ID, SYNCINFO_ID,
VERSIONING_DOCID, VERSIONING_DOCID,
type EntryVersionInfo,
type AnyEntry, type AnyEntry,
type EntryDoc, type EntryDoc,
type EntryLeaf, type EntryLeaf,
@@ -9,11 +9,6 @@ import {
type MetaEntry, type MetaEntry,
type ObsidianLiveSyncSettings, type ObsidianLiveSyncSettings,
} from "@vrtmrz/livesync-commonlib/compat/common/types"; } from "@vrtmrz/livesync-commonlib/compat/common/types";
import {
assessRemoteFeatureDocument,
describeRemoteFeatureRejection,
REMOTE_FEATURE_GENERATION,
} from "@vrtmrz/livesync-commonlib/replication";
import { isChunk } from "@vrtmrz/livesync-commonlib/compat/common/typeUtils"; import { isChunk } from "@vrtmrz/livesync-commonlib/compat/common/typeUtils";
import { import {
LOG_LEVEL_DEBUG, LOG_LEVEL_DEBUG,
@@ -69,10 +64,6 @@ interface ReplicateResultProcessorContext {
readonly services: ReplicateResultProcessorServices; readonly services: ReplicateResultProcessorServices;
} }
type ReplicateResultProcessorState = { type ReplicateResultProcessorState = {
databaseId?: string;
observedFeatures?: string[];
highestObservedVersion?: number;
invalidControlObserved?: boolean;
queued: PouchDB.Core.ExistingDocument<EntryDoc>[]; queued: PouchDB.Core.ExistingDocument<EntryDoc>[];
processing: PouchDB.Core.ExistingDocument<EntryDoc>[]; processing: PouchDB.Core.ExistingDocument<EntryDoc>[];
}; };
@@ -83,10 +74,6 @@ function shortenRev(rev: string | undefined): string {
if (!rev) return "undefined"; if (!rev) return "undefined";
return rev.length > 10 ? rev.substring(0, 10) : rev; return rev.length > 10 ? rev.substring(0, 10) : rev;
} }
function getPhysicalDatabaseId(database: PouchDB.Database<EntryDoc>): Promise<string | undefined> {
const identified = database as PouchDB.Database<EntryDoc> & { id?: () => Promise<string> };
return typeof identified.id === "function" ? identified.id() : Promise.resolve(undefined);
}
export class ReplicateResultProcessor { export class ReplicateResultProcessor {
private log(message: string, level: LOG_LEVEL = LOG_LEVEL_INFO) { private log(message: string, level: LOG_LEVEL = LOG_LEVEL_INFO) {
Logger(`[ReplicateResultProcessor] ${message}`, level); Logger(`[ReplicateResultProcessor] ${message}`, level);
@@ -129,89 +116,6 @@ export class ReplicateResultProcessor {
// If true, the processing queue processor bails the loop. // If true, the processing queue processor bails the loop.
private _suspended: boolean = false; private _suspended: boolean = false;
// A temporary lifecycle resume cannot make an unknown remote format safe to apply.
private _compatibilityBlocked = false;
private _assessingDatabase = false;
private _physicalDatabase?: PouchDB.Database<EntryDoc>;
private _observedFeatures = new Set<string>();
private _highestObservedVersion = 0;
private _invalidControlObserved = false;
public get isCompatibilityBlocked() {
return this._compatibilityBlocked;
}
private refreshPhysicalDatabase() {
const current = this.localDatabase.localDatabase;
if (current === this._physicalDatabase) return;
if (this._physicalDatabase) {
this._compatibilityBlocked = false;
this._assessingDatabase = true;
this._observedFeatures.clear();
this._highestObservedVersion = 0;
this._invalidControlObserved = false;
this._queuedChanges = [];
this._processingChanges = [];
this._restoreFromSnapshot = undefined;
this.updateProcessingActivity();
}
this._physicalDatabase = current;
}
private shouldStopApplication(sourceDatabase: PouchDB.Database<EntryDoc>) {
return (
this._compatibilityBlocked || this._assessingDatabase || sourceDatabase !== this.localDatabase.localDatabase
);
}
private blockForIncompatibleVersion(document: unknown, recordObservation = true) {
const assessment = assessRemoteFeatureDocument(document);
const hadObservedVersion = this._highestObservedVersion > 0;
let newFeaturesAdded = false;
if (recordObservation) {
let changed = false;
if (assessment.status === "invalid-control") {
changed = !this._invalidControlObserved;
this._invalidControlObserved = true;
} else {
const version = (document as EntryVersionInfo).version;
if (version > this._highestObservedVersion) {
this._highestObservedVersion = version;
changed = true;
}
if (assessment.status === "supported" || assessment.status === "unknown-features") {
for (const feature of assessment.status === "supported"
? assessment.usedFeatures
: ((document as EntryVersionInfo).used_features ?? [])) {
if (this._observedFeatures.has(feature)) continue;
this._observedFeatures.add(feature);
changed = true;
newFeaturesAdded = true;
}
}
}
if (changed) this.triggerTakeSnapshot();
}
if (assessment.status === "supported" || assessment.status === "older-generation") {
// A live writer must recheck the shared Tweak policy after another
// client starts using a newly declared representation.
if (
assessment.status === "supported" &&
hadObservedVersion &&
newFeaturesAdded &&
!this._assessingDatabase
) {
this.context.requestActiveReplicatorRetirement();
}
return;
}
if (this._compatibilityBlocked) return;
this._compatibilityBlocked = true;
this.updateProcessingActivity();
this.log(describeRemoteFeatureRejection(assessment), LOG_LEVEL_NOTICE);
this.context.requestActiveReplicatorRetirement();
}
/** /**
* Whether the application accepts replicated documents being applied. * Whether the application accepts replicated documents being applied.
* *
@@ -232,8 +136,6 @@ export class ReplicateResultProcessor {
public get isSuspended() { public get isSuspended() {
return ( return (
this._suspended || this._suspended ||
this._compatibilityBlocked ||
this._assessingDatabase ||
!this.acceptsResultApplication || !this.acceptsResultApplication ||
this.context.currentSettings().suspendParseReplicationResult || this.context.currentSettings().suspendParseReplicationResult ||
this.services.appLifecycle.isSuspended() this.services.appLifecycle.isSuspended()
@@ -244,38 +146,17 @@ export class ReplicateResultProcessor {
* Take a snapshot of the current processing state. * Take a snapshot of the current processing state.
* This snapshot is stored in the KV database for recovery on restart. * This snapshot is stored in the KV database for recovery on restart.
*/ */
private _snapshotWriter: Promise<void> = Promise.resolve(); protected async _takeSnapshot() {
const snapshot = {
protected _takeSnapshot(): Promise<void> { queued: this._queuedChanges.slice(),
// A blocked-batch flush must follow any earlier throttled write, or an processing: this._processingChanges.slice(),
// older snapshot could replace the queue after the replication callback. } satisfies ReplicateResultProcessorState;
const write = this._snapshotWriter await this.context.getKeyValueDB().set(KV_KEY_REPLICATION_RESULT_PROCESSOR_SNAPSHOT, snapshot);
.catch((): void => undefined) this.log(
.then(async () => { `Snapshot taken. Queued: ${snapshot.queued.length}, Processing: ${snapshot.processing.length}`,
const physicalDatabase = this.localDatabase.localDatabase; LOG_LEVEL_DEBUG
const databaseId = await getPhysicalDatabaseId(physicalDatabase); );
if (physicalDatabase !== this.localDatabase.localDatabase) return; this.reportStatus();
const snapshot = {
...(databaseId ? { databaseId } : {}),
observedFeatures: [...this._observedFeatures],
highestObservedVersion: this._highestObservedVersion,
invalidControlObserved: this._invalidControlObserved,
queued: this._queuedChanges.slice(),
processing: this._processingChanges.slice(),
} satisfies ReplicateResultProcessorState;
await this.context.getKeyValueDB().set(KV_KEY_REPLICATION_RESULT_PROCESSOR_SNAPSHOT, snapshot);
this.log(
`Snapshot taken. Queued: ${snapshot.queued.length}, Processing: ${snapshot.processing.length}`,
LOG_LEVEL_DEBUG
);
this.reportStatus();
});
this._snapshotWriter = write;
return write;
}
public async persistBlockedSnapshot(): Promise<void> {
if (this._compatibilityBlocked) await this._takeSnapshot();
} }
/** /**
* Trigger taking a snapshot. * Trigger taking a snapshot.
@@ -292,42 +173,10 @@ export class ReplicateResultProcessor {
* Restore from snapshot. * Restore from snapshot.
*/ */
public async restoreFromSnapshot() { public async restoreFromSnapshot() {
const physicalDatabase = this.localDatabase.localDatabase;
// Replication may have checkpointed a version document before its accompanying
// file changes reached the Vault. Assess the persisted requirement first.
let versionInfo: unknown;
try {
versionInfo = await this.localDatabase.getRaw(VERSIONING_DOCID);
} catch (error) {
if (!isNotFoundError(error)) throw error;
}
if (physicalDatabase !== this.localDatabase.localDatabase) return;
const snapshot = await this.context const snapshot = await this.context
.getKeyValueDB() .getKeyValueDB()
.get<ReplicateResultProcessorState>(KV_KEY_REPLICATION_RESULT_PROCESSOR_SNAPSHOT); .get<ReplicateResultProcessorState>(KV_KEY_REPLICATION_RESULT_PROCESSOR_SNAPSHOT);
if (physicalDatabase !== this.localDatabase.localDatabase) return; if (snapshot) {
const databaseId = await getPhysicalDatabaseId(physicalDatabase);
if (snapshot && (!snapshot.databaseId || !databaseId || snapshot.databaseId === databaseId)) {
for (const feature of snapshot.observedFeatures ?? []) this._observedFeatures.add(feature);
this._highestObservedVersion = Math.max(this._highestObservedVersion, snapshot.highestObservedVersion ?? 0);
this._invalidControlObserved ||= snapshot.invalidControlObserved === true;
}
if (versionInfo !== undefined) this.blockForIncompatibleVersion(versionInfo);
if (this._invalidControlObserved) this.blockForIncompatibleVersion(null, false);
if (this._highestObservedVersion > 0) {
this.blockForIncompatibleVersion(
{
_id: VERSIONING_DOCID,
type: "versioninfo",
version: this._highestObservedVersion,
...(this._highestObservedVersion >= REMOTE_FEATURE_GENERATION
? { used_features: [...this._observedFeatures] }
: {}),
},
false
);
}
if (snapshot && (!snapshot.databaseId || !databaseId || snapshot.databaseId === databaseId)) {
// Restoring the snapshot re-runs processing for both queued and processing items. // Restoring the snapshot re-runs processing for both queued and processing items.
const newQueue = [...snapshot.processing, ...snapshot.queued, ...this._queuedChanges]; const newQueue = [...snapshot.processing, ...snapshot.queued, ...this._queuedChanges];
this._queuedChanges = []; this._queuedChanges = [];
@@ -338,9 +187,6 @@ export class ReplicateResultProcessor {
); );
// await this._takeSnapshot(); // await this._takeSnapshot();
} }
this._assessingDatabase = false;
this.updateProcessingActivity();
this.triggerProcessQueue();
} }
private _restoreFromSnapshot: Promise<void> | undefined = undefined; private _restoreFromSnapshot: Promise<void> | undefined = undefined;
@@ -350,9 +196,7 @@ export class ReplicateResultProcessor {
* @returns Promise that resolves when restoration is complete. * @returns Promise that resolves when restoration is complete.
*/ */
public restoreFromSnapshotOnce() { public restoreFromSnapshotOnce() {
this.refreshPhysicalDatabase();
if (!this._restoreFromSnapshot) { if (!this._restoreFromSnapshot) {
this._assessingDatabase = true;
this._restoreFromSnapshot = this.restoreFromSnapshot(); this._restoreFromSnapshot = this.restoreFromSnapshot();
} }
return this._restoreFromSnapshot; return this._restoreFromSnapshot;
@@ -386,17 +230,7 @@ export class ReplicateResultProcessor {
* @param changes Changes to enqueue * @param changes Changes to enqueue
*/ */
public enqueueAll(changes: PouchDB.Core.ExistingDocument<EntryDoc>[], sourceDatabase?: PouchDB.Database<EntryDoc>) { public enqueueAll(changes: PouchDB.Core.ExistingDocument<EntryDoc>[]) {
if (sourceDatabase && sourceDatabase !== this.localDatabase.localDatabase) return;
const previousPhysicalDatabase = this._physicalDatabase;
this.refreshPhysicalDatabase();
if (previousPhysicalDatabase && previousPhysicalDatabase !== this._physicalDatabase) {
fireAndForget(() => this.restoreFromSnapshotOnce());
}
// Inspect every control document before a note in this batch can start applying.
for (const change of changes) {
if (change?._id === VERSIONING_DOCID) this.blockForIncompatibleVersion(change);
}
for (const change of changes) { for (const change of changes) {
// Check if the change is not a document change (e.g., chunk, versioninfo, syncinfo), and processed it directly. // Check if the change is not a document change (e.g., chunk, versioninfo, syncinfo), and processed it directly.
const isProcessed = this.processIfNonDocumentChange(change); const isProcessed = this.processIfNonDocumentChange(change);
@@ -441,8 +275,17 @@ export class ReplicateResultProcessor {
this.log(`Processed chunk: ${shortenId(change._id)}`, LOG_LEVEL_DEBUG); this.log(`Processed chunk: ${shortenId(change._id)}`, LOG_LEVEL_DEBUG);
return true; return true;
} }
if (change._id === VERSIONING_DOCID) { if (change._id === VERSIONING_DOCID || change.type === "versioninfo") {
this.log(`Version info document received: ${change._id}`, LOG_LEVEL_VERBOSE); this.log(`Version info document received: ${change._id}`, LOG_LEVEL_VERBOSE);
const assessment = assessRemoteFeatureDocument(change);
if (assessment.status !== "supported" && assessment.status !== "older-generation") {
// Fence and retire the active publication through its owner.
this.context.requestActiveReplicatorRetirement();
this.log(
`${describeRemoteFeatureRejection(assessment)} Update Self-hosted LiveSync before synchronising.`,
LOG_LEVEL_NOTICE
);
}
return true; return true;
} }
if ( if (
@@ -568,12 +411,11 @@ export class ReplicateResultProcessor {
// (per-document serialisation caps concurrency). // (per-document serialisation caps concurrency).
const releaser = await this._semaphore.acquire(); const releaser = await this._semaphore.acquire();
releaser(); releaser();
if (this.isSuspended) break;
// Dequeue the next change // Dequeue the next change
const doc = this._queuedChanges.shift(); const doc = this._queuedChanges.shift();
if (doc) { if (doc) {
this._processingChanges.push(doc); this._processingChanges.push(doc);
void this.parseDocumentChange(doc, this.localDatabase.localDatabase); void this.parseDocumentChange(doc);
} }
// Take snapshot (to be restored on next startup if needed) // Take snapshot (to be restored on next startup if needed)
this.triggerTakeSnapshot(); this.triggerTakeSnapshot();
@@ -589,12 +431,8 @@ export class ReplicateResultProcessor {
* @param change * @param change
* @returns * @returns
*/ */
async parseDocumentChange( async parseDocumentChange(change: PouchDB.Core.ExistingDocument<EntryDoc>) {
change: PouchDB.Core.ExistingDocument<EntryDoc>,
sourceDatabase: PouchDB.Database<EntryDoc> = this.localDatabase.localDatabase
) {
try { try {
if (this.shouldStopApplication(sourceDatabase)) return;
if (isAnyNote(change)) { if (isAnyNote(change)) {
const docMtime = change.mtime ?? 0; const docMtime = change.mtime ?? 0;
const maxMTime = this.context.currentSettings().maxMTimeForReflectEvents; const maxMTime = this.context.currentSettings().maxMTimeForReflectEvents;
@@ -611,7 +449,6 @@ export class ReplicateResultProcessor {
} }
// If the document is a virtual document, process it in the virtual document processor. // If the document is a virtual document, process it in the virtual document processor.
if (await this.services.replication.processVirtualDocument(change)) return; if (await this.services.replication.processVirtualDocument(change)) return;
if (this.shouldStopApplication(sourceDatabase)) return;
// If the document is version info, check compatibility and return. // If the document is version info, check compatibility and return.
if (isAnyNote(change)) { if (isAnyNote(change)) {
const docPath = this.getPath(change); const docPath = this.getPath(change);
@@ -619,7 +456,6 @@ export class ReplicateResultProcessor {
this.log(`Skipped: ${docPath}`, LOG_LEVEL_VERBOSE); this.log(`Skipped: ${docPath}`, LOG_LEVEL_VERBOSE);
return; return;
} }
if (this.shouldStopApplication(sourceDatabase)) return;
const size = change.size; const size = change.size;
// Note that this size check depends size that in metadata, not the actual content size. // Note that this size check depends size that in metadata, not the actual content size.
if (this.services.vault.isFileSizeTooLarge(size)) { if (this.services.vault.isFileSizeTooLarge(size)) {
@@ -629,20 +465,11 @@ export class ReplicateResultProcessor {
); );
return; return;
} }
return await this.applyToDatabase(change, sourceDatabase); return await this.applyToDatabase(change);
} }
this.log(`Skipped unexpected non-note document: ${change._id}`, LOG_LEVEL_INFO); this.log(`Skipped unexpected non-note document: ${change._id}`, LOG_LEVEL_INFO);
return; return;
} finally { } finally {
// An in-flight parse may have started before the control document arrived.
// Retain it even if a later asynchronous boundary stopped application.
if (
this._compatibilityBlocked &&
sourceDatabase === this.localDatabase.localDatabase &&
!this._queuedChanges.includes(change)
) {
this._queuedChanges.push(change);
}
// Remove from processing queue // Remove from processing queue
this._processingChanges = this._processingChanges.filter((e) => e !== change); this._processingChanges = this._processingChanges.filter((e) => e !== change);
try { try {
@@ -662,16 +489,12 @@ export class ReplicateResultProcessor {
} }
// Phase 2: apply the document to database // Phase 2: apply the document to database
protected applyToDatabase( protected applyToDatabase(doc: PouchDB.Core.ExistingDocument<AnyEntry>) {
doc: PouchDB.Core.ExistingDocument<AnyEntry>,
sourceDatabase: PouchDB.Database<EntryDoc> = this.localDatabase.localDatabase
) {
return this.withCounting(async () => { return this.withCounting(async () => {
let releaser: Awaited<ReturnType<typeof this._semaphore.acquire>> | undefined = undefined; let releaser: Awaited<ReturnType<typeof this._semaphore.acquire>> | undefined = undefined;
try { try {
releaser = await this._semaphore.acquire(); releaser = await this._semaphore.acquire();
if (this.shouldStopApplication(sourceDatabase)) return; await this._applyToDatabase(doc);
await this._applyToDatabase(doc, sourceDatabase);
} catch (e) { } catch (e) {
this.log(`Error while processing replication result`, LOG_LEVEL_NOTICE); this.log(`Error while processing replication result`, LOG_LEVEL_NOTICE);
this.logError(e); this.logError(e);
@@ -685,16 +508,12 @@ export class ReplicateResultProcessor {
} }
// Phase 2.1: process the document and apply to storage // Phase 2.1: process the document and apply to storage
// This function is serialized per document to avoid race-condition for the same document. // This function is serialized per document to avoid race-condition for the same document.
private _applyToDatabase( private _applyToDatabase(doc_: PouchDB.Core.ExistingDocument<AnyEntry>) {
doc_: PouchDB.Core.ExistingDocument<AnyEntry>,
sourceDatabase: PouchDB.Database<EntryDoc>
) {
const dbDoc = doc_ as LoadedEntry; // It has no `data` const dbDoc = doc_ as LoadedEntry; // It has no `data`
const path = this.getPath(dbDoc); const path = this.getPath(dbDoc);
return serialized(`replication-process:${dbDoc._id}`, async () => { return serialized(`replication-process:${dbDoc._id}`, async () => {
const docNote = `${path} (${shortenId(dbDoc._id)}, ${shortenRev(dbDoc._rev)})`; const docNote = `${path} (${shortenId(dbDoc._id)}, ${shortenRev(dbDoc._rev)})`;
const isRequired = await this.checkIsChangeRequiredForDatabaseProcessing(dbDoc); const isRequired = await this.checkIsChangeRequiredForDatabaseProcessing(dbDoc);
if (this.shouldStopApplication(sourceDatabase)) return;
if (!isRequired) { if (!isRequired) {
this.log(`Skipped (Not latest): ${docNote}`, LOG_LEVEL_VERBOSE); this.log(`Skipped (Not latest): ${docNote}`, LOG_LEVEL_VERBOSE);
return; return;
@@ -708,7 +527,6 @@ export class ReplicateResultProcessor {
const doc = isDeleted const doc = isDeleted
? { ...dbDoc, data: "" } ? { ...dbDoc, data: "" }
: await this.localDatabase.getDBEntryFromMeta({ ...dbDoc }, false, true); : await this.localDatabase.getDBEntryFromMeta({ ...dbDoc }, false, true);
if (this.shouldStopApplication(sourceDatabase)) return;
if (!doc) { if (!doc) {
// Failed to gather content // Failed to gather content
this.log(`Failed to gather content of ${docNote}`, LOG_LEVEL_NOTICE); this.log(`Failed to gather content of ${docNote}`, LOG_LEVEL_NOTICE);
@@ -719,10 +537,9 @@ export class ReplicateResultProcessor {
// Already processed // Already processed
this.log(`Processed by other processor: ${docNote}`, LOG_LEVEL_DEBUG); this.log(`Processed by other processor: ${docNote}`, LOG_LEVEL_DEBUG);
} else if (this.services.vault.isValidPath(this.getPath(doc))) { } else if (this.services.vault.isValidPath(this.getPath(doc))) {
if (this.shouldStopApplication(sourceDatabase)) return;
// Apply to storage if the path is valid // Apply to storage if the path is valid
try { try {
const reflected = await this.applyToStorage(doc as MetaEntry, sourceDatabase); const reflected = await this.applyToStorage(doc as MetaEntry);
if (!reflected) { if (!reflected) {
this.reportVaultReflectionFailure(doc as MetaEntry); this.reportVaultReflectionFailure(doc as MetaEntry);
return; return;
@@ -743,15 +560,9 @@ export class ReplicateResultProcessor {
* @param entry * @param entry
* @returns * @returns
*/ */
protected applyToStorage( protected applyToStorage(entry: MetaEntry) {
entry: MetaEntry,
sourceDatabase: PouchDB.Database<EntryDoc> = this.localDatabase.localDatabase
) {
return this.withCounting( return this.withCounting(
() => () => this.services.replication.processSynchroniseResult(entry),
this.shouldStopApplication(sourceDatabase)
? Promise.resolve(false)
: this.services.replication.processSynchroniseResult(entry),
this.services.replication.storageApplyingCount this.services.replication.storageApplyingCount
); );
} }
@@ -120,6 +120,34 @@ function setup(options: SetupOptions = {}) {
} }
describe("ReplicateResultProcessor", () => { describe("ReplicateResultProcessor", () => {
it("does not add a permanent application block when snapshot recovery fails", async () => {
const { processor } = setup({
getSnapshot: async () => {
throw new Error("KV unavailable");
},
});
await expect(processor.restoreFromSnapshotOnce()).rejects.toThrow("KV unavailable");
expect(processor.isSuspended).toBe(false);
});
it("restores pending notes without retaining a past feature rejection in KV", async () => {
const { processor, processSynchroniseResult, onCloseActiveReplication } = setup({
databaseId: "same-database",
getSnapshot: async () => ({
databaseId: "same-database",
invalidControlObserved: true,
observedFeatures: ["future-format-v7"],
observedGeneration: 14,
queued: [note("recovered-note")],
processing: [],
}),
});
await processor.restoreFromSnapshotOnce();
expect(processor.isSuspended).toBe(false);
await vi.waitFor(() => expect(processSynchroniseResult).toHaveBeenCalledOnce());
expect(onCloseActiveReplication).not.toHaveBeenCalled();
});
it.each([ it.each([
["Windows", isValidFilenameInWidows], ["Windows", isValidFilenameInWidows],
["Android", isValidFilenameInAndroid], ["Android", isValidFilenameInAndroid],
@@ -249,30 +277,6 @@ describe("ReplicateResultProcessor", () => {
expect(onCloseActiveReplication).not.toHaveBeenCalled(); expect(onCloseActiveReplication).not.toHaveBeenCalled();
}); });
it.each(["first", "last"] as const)(
"holds an entire received batch when an unknown feature is %s",
async (position) => {
const { onCloseActiveReplication, processor, processSynchroniseResult } = setup();
const versionInfo = {
_id: VERSIONING_DOCID,
_rev: "2-unknown",
type: "versioninfo",
version: REMOTE_FEATURE_GENERATION,
used_features: ["future-format-v7"],
} as PouchDB.Core.ExistingDocument<EntryDoc>;
const document = note("pending-after-compatibility-change");
processor.enqueueAll(position === "first" ? [versionInfo, document] : [document, versionInfo]);
await vi.waitFor(() => expect(onCloseActiveReplication).toHaveBeenCalledOnce());
await Promise.resolve();
expect(processor.isSuspended).toBe(true);
expect(processSynchroniseResult).not.toHaveBeenCalled();
processor.resume();
await Promise.resolve();
expect(processSynchroniseResult).not.toHaveBeenCalled();
}
);
it("continues when a newly received feature is supported", async () => { it("continues when a newly received feature is supported", async () => {
const { onCloseActiveReplication, processor, processSynchroniseResult } = setup(); const { onCloseActiveReplication, processor, processSynchroniseResult } = setup();
const versionInfo = { const versionInfo = {
@@ -289,149 +293,23 @@ describe("ReplicateResultProcessor", () => {
expect(onCloseActiveReplication).not.toHaveBeenCalled(); expect(onCloseActiveReplication).not.toHaveBeenCalled();
}); });
it("rechecks shared writer settings when a supported feature is added to an active database", async () => { it("reports unknown feature identifiers and requests Replicator retirement", () => {
const { onCloseActiveReplication, processor } = setup({ const logger = vi.fn();
localVersionInfo: { setGlobalLogFunction(logger);
_id: VERSIONING_DOCID,
type: "versioninfo",
version: REMOTE_FEATURE_GENERATION,
used_features: [],
},
});
await processor.restoreFromSnapshotOnce();
const versionInfo = {
_id: VERSIONING_DOCID,
_rev: "2-supported",
type: "versioninfo",
version: REMOTE_FEATURE_GENERATION,
used_features: [ENCRYPTED_INTERNAL_METADATA_FEATURE],
} as PouchDB.Core.ExistingDocument<EntryDoc>;
processor.enqueueAll([versionInfo]);
processor.enqueueAll([{ ...versionInfo, _rev: "3-same-features" }]);
expect(onCloseActiveReplication).toHaveBeenCalledOnce();
expect(processor.isSuspended).toBe(false);
});
it("retains an in-flight note when a feature change arrives during an asynchronous check", async () => {
const targetCheck = promiseWithResolvers<boolean>();
const { isTargetFile, onCloseActiveReplication, processor, processSynchroniseResult } = setup({
isTargetFile: async () => targetCheck.promise,
});
const pending = note("in-flight");
processor.enqueueAll([pending]);
await vi.waitFor(() => expect(isTargetFile).toHaveBeenCalledOnce());
processor.enqueueAll([
{
_id: VERSIONING_DOCID,
_rev: "2-unknown",
type: "versioninfo",
version: REMOTE_FEATURE_GENERATION,
used_features: ["future-format-v7"],
} as PouchDB.Core.ExistingDocument<EntryDoc>,
]);
targetCheck.resolve(true);
await vi.waitFor(() => expect(processor["_processingChanges"]).toHaveLength(0));
expect(processor["_queuedChanges"]).toContain(pending);
expect(processSynchroniseResult).not.toHaveBeenCalled();
expect(onCloseActiveReplication).toHaveBeenCalledOnce();
});
it("restores a checkpointed note only after assessing the persisted version document", async () => {
const pending = note("checkpointed");
const { onCloseActiveReplication, processor, processSynchroniseResult } = setup({
localVersionInfo: {
_id: VERSIONING_DOCID,
type: "versioninfo",
version: REMOTE_FEATURE_GENERATION,
used_features: ["future-format-v7"],
},
getSnapshot: async () => ({ processing: [pending], queued: [] }),
});
await processor.restoreFromSnapshotOnce();
expect(processor.isSuspended).toBe(true);
expect(processor["_queuedChanges"]).toContain(pending);
expect(processSynchroniseResult).not.toHaveBeenCalled();
expect(onCloseActiveReplication).toHaveBeenCalledOnce();
});
it("retains an observed unknown feature when the local version list is later shortened", async () => {
let snapshot: unknown;
const first = setup({
databaseId: "same-database",
setSnapshot: async (_key, value) => {
snapshot = value;
},
});
const pending = note("checkpointed-after-list-change");
first.processor.enqueueAll([
{
_id: VERSIONING_DOCID,
_rev: "2-unknown",
type: "versioninfo",
version: REMOTE_FEATURE_GENERATION,
used_features: ["future-format-v7"],
} as PouchDB.Core.ExistingDocument<EntryDoc>,
pending,
]);
await vi.waitFor(() => expect(snapshot).toBeDefined());
const resumed = setup({
databaseId: "same-database",
localVersionInfo: {
_id: VERSIONING_DOCID,
type: "versioninfo",
version: REMOTE_FEATURE_GENERATION,
used_features: [],
},
getSnapshot: async () => snapshot,
});
await resumed.processor.restoreFromSnapshotOnce();
expect(resumed.processor.isSuspended).toBe(true);
expect(resumed.processor["_queuedChanges"]).toContain(pending);
expect(resumed.processSynchroniseResult).not.toHaveBeenCalled();
expect(resumed.onCloseActiveReplication).toHaveBeenCalledOnce();
});
it("does not restore pending work belonging to a replaced physical database", async () => {
const pending = note("old-database");
const { processor, processSynchroniseResult } = setup({
databaseId: "replacement-database",
getSnapshot: async () => ({ databaseId: "retired-database", processing: [pending], queued: [] }),
});
await processor.restoreFromSnapshotOnce();
expect(processor["_queuedChanges"]).toHaveLength(0);
expect(processSynchroniseResult).not.toHaveBeenCalled();
expect(processor.isSuspended).toBe(false);
});
it("reports unknown identifiers and retires only once for repeated notifications", () => {
const log = vi.fn((_message: unknown, _level?: number) => undefined);
setGlobalLogFunction(log);
try { try {
const { onCloseActiveReplication, processor } = setup(); const { processor, onCloseActiveReplication } = setup();
const versionInfo = { processor.enqueueAll([
_id: VERSIONING_DOCID, {
_rev: "2-unknown", _id: VERSIONING_DOCID,
type: "versioninfo", _rev: "1-unknown",
version: REMOTE_FEATURE_GENERATION, type: "versioninfo",
used_features: ["future-format-v7"], version: REMOTE_FEATURE_GENERATION,
} as PouchDB.Core.ExistingDocument<EntryDoc>; used_features: ["future-format-v7"],
} as PouchDB.Core.ExistingDocument<EntryDoc>,
processor.enqueueAll([versionInfo]); ]);
processor.enqueueAll([versionInfo]);
expect(onCloseActiveReplication).toHaveBeenCalledOnce(); expect(onCloseActiveReplication).toHaveBeenCalledOnce();
expect(log).toHaveBeenCalledWith( expect(logger).toHaveBeenCalledWith(
"[ReplicateResultProcessor] Unknown features are in use: future-format-v7", expect.stringContaining("future-format-v7"),
LOG_LEVEL_NOTICE, LOG_LEVEL_NOTICE,
undefined undefined
); );
@@ -440,43 +318,6 @@ describe("ReplicateResultProcessor", () => {
} }
}); });
it("ignores a late feature notification from a replaced physical database", () => {
const { onCloseActiveReplication, processor } = setup();
const oldPhysicalDatabase = {} as PouchDB.Database<EntryDoc>;
const versionInfo = {
_id: VERSIONING_DOCID,
_rev: "2-unknown",
type: "versioninfo",
version: REMOTE_FEATURE_GENERATION,
used_features: ["future-format-v7"],
} as PouchDB.Core.ExistingDocument<EntryDoc>;
processor.enqueueAll([versionInfo], oldPhysicalDatabase);
expect(processor.isSuspended).toBe(false);
expect(onCloseActiveReplication).not.toHaveBeenCalled();
});
it("reassesses compatibility for a replacement local database", async () => {
const { localDatabase, onCloseActiveReplication, processor } = setup();
processor.enqueueAll([
{
_id: VERSIONING_DOCID,
_rev: "2-unknown",
type: "versioninfo",
version: REMOTE_FEATURE_GENERATION,
used_features: ["future-format-v7"],
} as PouchDB.Core.ExistingDocument<EntryDoc>,
]);
expect(processor.isSuspended).toBe(true);
localDatabase.localDatabase = {} as PouchDB.Database<EntryDoc>;
await processor.restoreFromSnapshotOnce();
expect(processor.isSuspended).toBe(false);
expect(onCloseActiveReplication).toHaveBeenCalledOnce();
});
it("scans normal-file metadata without loading chunk documents and requeues it", async () => { it("scans normal-file metadata without loading chunk documents and requeues it", async () => {
const documents = [ const documents = [
{ _id: "first", _rev: "1-a", type: "plain", path: "first.md" }, { _id: "first", _rev: "1-a", type: "plain", path: "first.md" },
+6 -8
View File
@@ -98,21 +98,19 @@ export function useReplicationFeature<TContext extends ServiceContext, TCommands
clearHandlers(); clearHandlers();
return Promise.resolve(true); return Promise.resolve(true);
}); });
services.databaseEvents.onDatabaseInitialised.addHandler(async () => { services.databaseEvents.onDatabaseInitialised.addHandler(() => {
await resultProcessor.restoreFromSnapshotOnce(); fireAndForget(() => resultProcessor.restoreFromSnapshotOnce());
return true; return Promise.resolve(true);
}); });
services.appLifecycle.onSettingLoaded.addHandler(initialiseAutomaticReplicationTriggers); services.appLifecycle.onSettingLoaded.addHandler(initialiseAutomaticReplicationTriggers);
services.replication.parseSynchroniseResult.addHandler(async (documents, sourceDatabase) => { services.replication.parseSynchroniseResult.addHandler((documents) => {
resultProcessor.enqueueAll(documents, sourceDatabase); resultProcessor.enqueueAll(documents);
await resultProcessor.persistBlockedSnapshot(); return Promise.resolve(true);
return true;
}); });
services.replication.onBeforeReplicate.addHandler(onlinePreflight, 10); services.replication.onBeforeReplicate.addHandler(onlinePreflight, 10);
services.replication.onPrepareCentralRemoteReplication.addHandler(securitySeedPreflight); services.replication.onPrepareCentralRemoteReplication.addHandler(securitySeedPreflight);
services.replication.onBeforeReplicate.addHandler(async () => { services.replication.onBeforeReplicate.addHandler(async () => {
await resultProcessor.restoreFromSnapshotOnce(); await resultProcessor.restoreFromSnapshotOnce();
if (resultProcessor.isCompatibilityBlocked) return false;
unresolvedErrorManager.clearErrors(); unresolvedErrorManager.clearErrors();
return true; return true;
}, 100); }, 100);
@@ -22,7 +22,9 @@ type SetupOptions = {
function setup(options: SetupOptions = {}) { function setup(options: SetupOptions = {}) {
const defaultLocalDatabase = { const defaultLocalDatabase = {
localDatabase: {}, localDatabase: {},
getRaw: vi.fn(async () => { throw { status: 404 }; }), getRaw: vi.fn(async () => {
throw { status: 404 };
}),
}; };
const { const {
getLocalDatabase = () => defaultLocalDatabase, getLocalDatabase = () => defaultLocalDatabase,
@@ -183,59 +185,4 @@ describe("replication serviceFeature composition", () => {
retirement.resolve(true); retirement.resolve(true);
}); });
it("persists a blocked batch before its replication callback settles", async () => {
const writeFinished = promiseWithResolvers<void>();
const writes: Array<{ queued: Array<{ _id: string }> }> = [];
const { parseHandler } = setup({
keyValueDB: {
kvDB: {
get: vi.fn(async () => undefined),
set: vi.fn(async (_key, value) => {
writes.push(value as { queued: Array<{ _id: string }> });
await writeFinished.promise;
}),
},
},
});
const pending = {
_id: "checkpointed-note",
_rev: "1-test",
type: "plain",
path: "checkpointed-note.md",
} as PouchDB.Core.ExistingDocument<EntryDoc>;
const unknownVersion = {
_id: VERSIONING_DOCID,
_rev: "2-unknown",
type: "versioninfo",
version: REMOTE_FEATURE_GENERATION,
used_features: ["future-format-v7"],
} as PouchDB.Core.ExistingDocument<EntryDoc>;
let callbackSettled = false;
const callback = parseHandler!([pending, unknownVersion]).then((result) => {
callbackSettled = true;
return result;
});
await vi.waitFor(() => expect(writes.length).toBeGreaterThan(0));
expect(callbackSettled).toBe(false);
writeFinished.resolve();
await expect(callback).resolves.toBe(true);
expect(writes[writes.length - 1]?.queued.map((entry) => entry._id)).toContain(pending._id);
});
it("refuses another replication after observing an unknown feature locally", async () => {
const { beforeReplicateHandlers, parseHandler } = setup();
const unknownVersion = {
_id: VERSIONING_DOCID,
_rev: "2-unknown",
type: "versioninfo",
version: REMOTE_FEATURE_GENERATION,
used_features: ["future-format-v7"],
} as PouchDB.Core.ExistingDocument<EntryDoc>;
await parseHandler!([unknownVersion]);
await expect(beforeReplicateHandlers.get(100)!(false)).resolves.toBe(false);
});
}); });
@@ -368,13 +368,20 @@ async function uploadHiddenFile(
return ids.has(entry.id) && entry.children.every((childId) => ids.has(childId)); return ids.has(entry.id) && entry.children.every((childId) => ids.has(childId));
}); });
const remoteEntry = await fetchCouchDbDocument(context.couchDb, context.dbName, entry.id); const remoteEntry = await fetchCouchDbDocument(context.couchDb, context.dbName, entry.id);
if (!remoteEntry.path?.startsWith("/\\:") || remoteEntry.children?.length !== 0 || if (
remoteEntry.ctime !== 0 || remoteEntry.mtime !== 0 || remoteEntry.size !== 0) { !remoteEntry.path?.startsWith("/\\:") ||
remoteEntry.children?.length !== 0 ||
remoteEntry.ctime !== 0 ||
remoteEntry.mtime !== 0 ||
remoteEntry.size !== 0
) {
throw new Error(`Hidden File Sync Metadata was not encrypted for ${entry.id}.`); throw new Error(`Hidden File Sync Metadata was not encrypted for ${entry.id}.`);
} }
const versionInfo = await fetchCouchDbDocument(context.couchDb, context.dbName, VERSIONING_DOCID); const versionInfo = await fetchCouchDbDocument(context.couchDb, context.dbName, VERSIONING_DOCID);
if (versionInfo.version !== REMOTE_FEATURE_GENERATION || if (
!(versionInfo.used_features as unknown[] | undefined)?.includes(ENCRYPTED_INTERNAL_METADATA_FEATURE)) { versionInfo.version !== REMOTE_FEATURE_GENERATION ||
!(versionInfo.used_features as unknown[] | undefined)?.includes(ENCRYPTED_INTERNAL_METADATA_FEATURE)
) {
throw new Error("The remote feature list does not declare encrypted internal Metadata."); throw new Error("The remote feature list does not declare encrypted internal Metadata.");
} }
return entry; return entry;
@@ -400,6 +407,29 @@ async function runCreateRoundTrip(
await writeVaultFile(vaultA.path, snippetPath, snippetContent); await writeVaultFile(vaultA.path, snippetPath, snippetContent);
let session = await startConfiguredSession(context, vaultA); let session = await startConfiguredSession(context, vaultA);
const entry = await uploadHiddenFile(context, session, snippetPath); const entry = await uploadHiddenFile(context, session, snippetPath);
await evalObsidianJson(
context.cliBinary,
[
"(async()=>{",
"const rebuilder=app.plugins.plugins['obsidian-livesync'].core.rebuilder;",
"const inform=rebuilder.informOptionalFeatures;",
"rebuilder.informOptionalFeatures=async()=>{};",
"try{await rebuilder.$rebuildRemote();}finally{rebuilder.informOptionalFeatures=inform;}",
"return JSON.stringify(true);",
"})()",
].join(""),
session.cliEnv
);
const rebuiltVersion = await fetchCouchDbDocument(context.couchDb, context.dbName, VERSIONING_DOCID);
const rebuiltEntry = await fetchCouchDbDocument(context.couchDb, context.dbName, entry.id);
if (
rebuiltVersion.version !== REMOTE_FEATURE_GENERATION ||
!(rebuiltVersion.used_features as unknown[] | undefined)?.includes(ENCRYPTED_INTERNAL_METADATA_FEATURE) ||
!rebuiltEntry.path?.startsWith("/\\:")
) {
throw new Error("Remote Rebuild did not declare and encrypt internal Metadata for its accepted writer.");
}
console.log("Remote Rebuild declared the feature and preserved encrypted Hidden File Sync Metadata.");
await session.app.stop(); await session.app.stop();
session = await startConfiguredSession(context, vaultB); session = await startConfiguredSession(context, vaultB);
@@ -588,11 +618,14 @@ async function runInitialisationNoticeGrouping(context: RunnerContext, vault: Te
await withObsidianPage(port, async (page) => { await withObsidianPage(port, async (page) => {
const deadline = Date.now() + timeoutMs; const deadline = Date.now() + timeoutMs;
while ((await page.locator(".notice:visible").count()) > 0 && Date.now() < deadline) { while ((await page.locator(".notice:visible").count()) > 0 && Date.now() < deadline) {
await page.locator(".notice:visible").first().click({ await page
force: true, .locator(".notice:visible")
position: { x: 2, y: 2 }, .first()
timeout: timeoutMs, .click({
}); force: true,
position: { x: 2, y: 2 },
timeout: timeoutMs,
});
} }
assertEqual( assertEqual(
await page.locator(".notice:visible").count(), await page.locator(".notice:visible").count(),
@@ -724,17 +757,15 @@ async function runInitialisationNoticeGrouping(context: RunnerContext, vault: Te
const result = await withObsidianPage(port, async (page) => { const result = await withObsidianPage(port, async (page) => {
await page.evaluate((stateKey) => { await page.evaluate((stateKey) => {
const state = (globalThis as unknown as Record< const state = (
string, globalThis as unknown as Record<string, { releasePreparation?: () => void } | undefined>
{ releasePreparation?: () => void } | undefined )[stateKey];
>)[stateKey];
state?.releasePreparation?.(); state?.releasePreparation?.();
}, hiddenFileInitialisationStateKey); }, hiddenFileInitialisationStateKey);
await page.waitForFunction( await page.waitForFunction(
(stateKey) => (stateKey) =>
(globalThis as unknown as Record<string, { reachedInitialisation?: boolean } | undefined>)[ (globalThis as unknown as Record<string, { reachedInitialisation?: boolean } | undefined>)[stateKey]
stateKey ?.reachedInitialisation === true,
]?.reachedInitialisation === true,
hiddenFileInitialisationStateKey, hiddenFileInitialisationStateKey,
{ timeout: timeoutMs } { timeout: timeoutMs }
); );
@@ -32,9 +32,6 @@ const unknownFeature = "future-format-v7";
type FeatureState = { type FeatureState = {
version: number | null; version: number | null;
features: string[]; features: string[];
snapshotFeatures: string[];
snapshotDatabaseId: string | null;
databaseId: string | null;
hasActiveReplicator: boolean; hasActiveReplicator: boolean;
}; };
@@ -46,14 +43,9 @@ async function readFeatureState(cliBinary: string, env: NodeJS.ProcessEnv): Prom
"const core=app.plugins.plugins['obsidian-livesync'].core;", "const core=app.plugins.plugins['obsidian-livesync'].core;",
`const id=${JSON.stringify(VERSIONING_DOCID)};`, `const id=${JSON.stringify(VERSIONING_DOCID)};`,
"const info=await core.localDatabase.getRaw(id).catch(()=>null);", "const info=await core.localDatabase.getRaw(id).catch(()=>null);",
"const snapshot=await core.kvDB.get('replicationResultProcessorSnapshot');",
"const databaseId=await core.localDatabase.localDatabase.id();",
"return JSON.stringify({", "return JSON.stringify({",
"version:typeof info?.version==='number'?info.version:null,", "version:typeof info?.version==='number'?info.version:null,",
"features:Array.isArray(info?.used_features)?info.used_features:[],", "features:Array.isArray(info?.used_features)?info.used_features:[],",
"snapshotFeatures:Array.isArray(snapshot?.observedFeatures)?snapshot.observedFeatures:[],",
"snapshotDatabaseId:snapshot?.databaseId??null,",
"databaseId,",
"hasActiveReplicator:!!core.services.replicator.getActiveReplicator(),", "hasActiveReplicator:!!core.services.replicator.getActiveReplicator(),",
"});", "});",
"})()", "})()",
@@ -153,12 +145,8 @@ async function main(): Promise<void> {
const observed = await waitForState( const observed = await waitForState(
cli.binary, cli.binary,
session.cliEnv, session.cliEnv,
(state) => (state) => state.version === 13 && state.features.includes(unknownFeature) && !state.hasActiveReplicator,
state.version === 13 && "the live feature change and Replicator retirement"
state.features.includes(unknownFeature) &&
state.snapshotFeatures.includes(unknownFeature) &&
!state.hasActiveReplicator,
"the live feature change, durable observation, and Replicator retirement"
); );
const replicated = await evalObsidianJson<boolean>( const replicated = await evalObsidianJson<boolean>(
@@ -171,25 +159,6 @@ async function main(): Promise<void> {
if (acceptedAfterStop !== acceptedContent) if (acceptedAfterStop !== acceptedContent)
throw new Error("Previously accepted Vault content changed on stop."); throw new Error("Previously accepted Vault content changed on stop.");
// Shorten both visible lists so only the snapshot can retain the observed requirement.
const shortenedRemoteVersion = await fetchCouchDbDocument(couchDb, dbName, VERSIONING_DOCID);
await putCouchDbDocument(couchDb, dbName, { ...shortenedRemoteVersion, used_features: [] });
const shortenedLocalFeatures = await evalObsidianJson<string[]>(
cli.binary,
[
"(async()=>{",
"const core=app.plugins.plugins['obsidian-livesync'].core;",
`const id=${JSON.stringify(VERSIONING_DOCID)};`,
"const info=await core.localDatabase.getRaw(id);",
"await core.localDatabase.putRaw({...info,used_features:[]});",
"return JSON.stringify((await core.localDatabase.getRaw(id)).used_features);",
"})()",
].join(""),
session.cliEnv
);
if (shortenedLocalFeatures.length !== 0)
throw new Error("The local control document did not shorten for the restart scenario.");
await session.app.stop(); await session.app.stop();
session = undefined; session = undefined;
session = await startObsidianLiveSyncSession({ binary, cliBinary: cli.binary, vault }); session = await startObsidianLiveSyncSession({ binary, cliBinary: cli.binary, vault });
@@ -197,36 +166,40 @@ async function main(): Promise<void> {
const afterRestart = await waitForState( const afterRestart = await waitForState(
cli.binary, cli.binary,
session.cliEnv, session.cliEnv,
(state) => (state) => state.version === 13 && state.features.includes(unknownFeature),
state.version === 13 && state.features.length === 0 && state.snapshotFeatures.includes(unknownFeature), "the remote feature requirement after restart"
"the retained unknown feature after restart"
); );
if (afterRestart.snapshotDatabaseId !== afterRestart.databaseId)
throw new Error("The retained feature observation belongs to a different physical database.");
const replicatedAfterRestart = await evalObsidianJson<boolean>( const replicatedAfterRestart = await evalObsidianJson<boolean>(
cli.binary, cli.binary,
"(async()=>JSON.stringify(!!(await app.plugins.plugins['obsidian-livesync'].core.services.replication.replicate(true))))()", "(async()=>JSON.stringify(!!(await app.plugins.plugins['obsidian-livesync'].core.services.replication.replicate(true))))()",
session.cliEnv session.cliEnv
); );
const continuousAfterRestart = await evalObsidianJson<{ status: string }>( const continuousAfterRestart = await evalObsidianJson<{ requestStatus: string; connected: boolean }>(
cli.binary, cli.binary,
[ [
"(async()=>{", "(async()=>{",
"const core=app.plugins.plugins['obsidian-livesync'].core;", "const core=app.plugins.plugins['obsidian-livesync'].core;",
"const replicator=core.services.replicator.getActiveReplicator();",
"const original=replicator.openContinuousReplication;",
"let completion;",
"replicator.openContinuousReplication=function(...args){completion=original.apply(this,args);return completion;};",
"try{",
"const result=await core.services.replication.startContinuous({trigger:'daemon',interaction:{kind:'forbidden'}});", "const result=await core.services.replication.startContinuous({trigger:'daemon',interaction:{kind:'forbidden'}});",
"return JSON.stringify(result);", "if(!completion) throw new Error('The continuous provider did not attempt its remote check');",
"return JSON.stringify({requestStatus:result.status,connected:await completion});",
"}finally{replicator.openContinuousReplication=original;}",
"})()", "})()",
].join(""), ].join(""),
session.cliEnv session.cliEnv
); );
if (replicatedAfterRestart || continuousAfterRestart.status !== "blocked") if (replicatedAfterRestart || continuousAfterRestart.connected !== false)
throw new Error( throw new Error(
`Replication resumed after restart despite a previously observed unknown feature: ${JSON.stringify({ afterRestart, replicatedAfterRestart, continuousAfterRestart })}` `Replication resumed after restart despite an unknown remote feature: ${JSON.stringify({ afterRestart, replicatedAfterRestart, continuousAfterRestart })}`
); );
if ((await readFile(fullPath, "utf-8")) !== acceptedContent) if ((await readFile(fullPath, "utf-8")) !== acceptedContent)
throw new Error("Previously accepted Vault content changed after restart."); throw new Error("Previously accepted Vault content changed after restart.");
console.log( console.log(
`Active feature change retired the Replicator; a shortened control document did not clear the block after restart: ${JSON.stringify({ observed, afterRestart })}` `Active feature change retired the Replicator; the remote declaration refused synchronisation after restart: ${JSON.stringify({ observed, afterRestart })}`
); );
} finally { } finally {
await session?.app.stop(); await session?.app.stop();
+1 -1
View File
@@ -17,7 +17,7 @@ Earlier releases remain available in the 1.0 release history, the 1.0 preview hi
#### New Feature #### New Feature
- Hidden File Sync and Customisation Sync can now encrypt their paths, times, sizes, and Chunk references in CouchDB Metadata when E2EE V2 and Property Encryption are enabled. The preference is enabled for new Vaults and remains off for existing configurations until selected. It protects future writes; protecting existing Metadata also requires a manual remote Rebuild after all devices have been updated. - Hidden File Sync and Customisation Sync can now encrypt their paths, times, sizes, and Chunk references in CouchDB Metadata when E2EE V2 and Property Encryption are enabled. The preference is enabled for new Vaults and remains off for existing configurations until selected. It protects future writes; protecting existing Metadata also requires a manual remote Rebuild after all devices have been updated.
- CouchDB records the features its data uses. Clients now stop synchronisation and pending file reflection when they encounter an unknown feature, and show its identifier so that the required update can be identified. - CouchDB records the features its data uses. Clients check these requirements before synchronisation and show any unknown feature identifiers. Receiving an unsupported requirement also stops the active replication.
## 1.0.32 ## 1.0.32