diff --git a/docs/design_docs/internal_metadata_encryption.md b/docs/design_docs/internal_metadata_encryption.md index 5ceaf218..1ec82835 100644 --- a/docs/design_docs/internal_metadata_encryption.md +++ b/docs/design_docs/internal_metadata_encryption.md @@ -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 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 - `parseSynchroniseResult` with the received documents. -2. The replication service feature passes them to - `ReplicateResultProcessor.enqueueAll`. -3. `processIfNonDocumentChange` recognises `type: versioninfo` and requests active - Replicator retirement when `version > VER`. -4. The owner closes admission, requests transfer cancellation, drains its work, - and closes the instance. The result callback does not wait for that transition. +Retain the existing received-version path through `parseSynchroniseResult`, +`enqueueAll`, and `processIfNonDocumentChange`. Replace its numeric comparison +with the shared assessment so unknown names at the same generation are also +reported. Known features, reordered lists, and ordinary revision updates do not +retire the connection. Unsupported or malformed control documents request +retirement through the existing Replicator owner and display the reason. +The callback must not await retirement of the operation which delivered it. -Retain this observation path, but use Commonlib's complete assessment of the -identified control document. A changed `used_features` list must be inspected -even when the numeric version is unchanged. A mere revision change, list reorder, -or duplicate known identifier does not require an interruption. - -Inspect the entire batch's control information before passing any file entries -to normal or optional processing. Recognise the fixed control-document ID and -validate its type and contents. A deleted or malformed control document is a -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. +This is an admission check and a best-effort stop for exceptional changes during +an active connection. It does not fence every queued file application, roll back +accepted writes, or guarantee an atomic change across live devices. Feature +changes are an infrequent administrative operation: update all devices first, +then enable the preference and use the recommended manual Rebuild. Rebuild +locks the remote using the existing workflow; changing this preference alone +does not lock it. The action to proceed without rebuilding explicitly reminds +the person to update every other device, including currently connected devices. ## Persistence and recovery boundaries -A CouchDB replication notification can arrive after the documents have entered -the local DB. Already-started network and filesystem operations may settle. -This feature does not promise rollback or atomic revocation of those operations. +Do not retain a second feature list, highest generation, or rejection flag in +KV storage. Do not add compatibility checks to pending-work snapshot recovery +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 -checkpoint may already include the documents which have not reached the Vault. -Do not drop those documents or depend on an ordinary reconnect to send them -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. +After updating clients, use normal reconnection and the existing Hatch +inspection or Fetch workflow if reconciliation is needed. This feature does not +repair unrelated KV inconsistencies or the existing readiness queue behaviour. Garbage Collection V3 is a beta manual operation which begins with an ordinary 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 automatic Rebuild, Fetch, or restart for this rule. -Add deterministic runtime tests for feature-only changes, version documents -first and last in a batch, unknown-name presentation, duplicate notifications, -queued and waiting reflection, stale database callbacks, restart, checkpointed -but unapplied documents, and cancellation without a circular wait. +Keep unit tests for known and unknown feature notifications, generic identifier +presentation, retirement without a circular wait, and the unchanged snapshot +behaviour after KV failure or obsolete snapshot fields. The previous batch +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 -incompatibility requests retirement, but the processor still applies a note in -the same batch when its host remains ready. A same-version document with an -unknown feature does not request retirement. These observations motivated the -new checks. - -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. +Use real Obsidian Hidden File Sync and Customisation Sync scenarios to inspect +raw CouchDB Metadata and restore content in another Vault. Check the admitted +writer on a locked remote, unknown-feature rejection before and during +replication, and remote-based rejection after restart. Retain the encrypted +CLI-to-Obsidian interoperability scenario. A future client upgrade that adds +support for an unknown feature is a separate validation boundary. Keep the primary-language settings and troubleshooting guides, the database-compatibility ADR, and Unreleased notes aligned with this behaviour. diff --git a/docs/settings.md b/docs/settings.md index d6622d91..3cff2855 100644 --- a/docs/settings.md +++ b/docs/settings.md @@ -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. -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 diff --git a/docs/troubleshooting.md b/docs/troubleshooting.md index b2a00a68..95489ba1 100644 --- a/docs/troubleshooting.md +++ b/docs/troubleshooting.md @@ -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 -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 diff --git a/src/modules/features/SettingDialogue/PaneRemoteConfig.ts b/src/modules/features/SettingDialogue/PaneRemoteConfig.ts index 3f93b00c..e193c986 100644 --- a/src/modules/features/SettingDialogue/PaneRemoteConfig.ts +++ b/src/modules/features/SettingDialogue/PaneRemoteConfig.ts @@ -5,7 +5,6 @@ import { DEFAULT_SETTINGS, LOG_LEVEL_NOTICE, type ObsidianLiveSyncSettings, - type EncryptionSettings, LOG_LEVEL_VERBOSE, } from "@vrtmrz/livesync-commonlib/compat/common/types"; import { Menu, type ButtonComponent } from "@/deps.ts"; @@ -32,13 +31,11 @@ import { import { ConnectionStringParser } 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 SetupRemoteE2EE from "@/modules/features/SetupWizard/dialogs/SetupRemoteE2EE.svelte"; import SetupRemoteCouchDB from "@/modules/features/SetupWizard/dialogs/SetupRemoteCouchDB.svelte"; import SetupRemoteBucket from "@/modules/features/SetupWizard/dialogs/SetupRemoteBucket.svelte"; import type { SetupRemoteCouchDBInitialData, SetupRemoteCouchDBResultType, - SetupRemoteE2EEResultType, } from "@/modules/features/SetupWizard/dialogs/setupDialogTypes.ts"; import { syncActivatedRemoteSettings } from "./remoteConfigBuffer.ts"; @@ -119,34 +116,15 @@ export function paneRemoteConfig( .onClick(async () => { const setupManager = this.core.getModule(SetupManager); const originalSettings = getSettingsFromEditingSettings(this.editingSettings); - const e2eeConf = await setupManager.dialogManager.openWithExplicitCancel< - SetupRemoteE2EEResultType, - EncryptionSettings - >(SetupRemoteE2EE, originalSettings); - 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; + const applied = await setupManager.onlyE2EEConfiguration(UserMode.Update, originalSettings); + if (applied) { + this.editingSettings.encryptInternalMetadata = + this.core.settings.encryptInternalMetadata; if (this.initialSettings) { - this.initialSettings.encryptInternalMetadata = e2eeConf.encryptInternalMetadata; + this.initialSettings.encryptInternalMetadata = + this.core.settings.encryptInternalMetadata; } this.requestUpdate(); - } else { - await setupManager.onConfirmApplySettingsFromWizard( - { ...originalSettings, ...e2eeConf }, - UserMode.Update - ); } updateE2EESummary(); }) diff --git a/src/modules/features/SettingDialogue/PaneRemoteConfig.unit.spec.ts b/src/modules/features/SettingDialogue/PaneRemoteConfig.unit.spec.ts index 8d259210..efa79c83 100644 --- a/src/modules/features/SettingDialogue/PaneRemoteConfig.unit.spec.ts +++ b/src/modules/features/SettingDialogue/PaneRemoteConfig.unit.spec.ts @@ -162,24 +162,15 @@ describe("paneRemoteConfig", () => { encryptInternalMetadata: false, remoteConfigurations: {}, }; - const applyPartial = vi.fn(async () => {}); - const onConfirmApplySettingsFromWizard = vi.fn(async () => {}); const setupManager = { - dialogManager: { - openWithExplicitCancel: vi.fn(async () => ({ - encrypt: true, - passphrase: "passphrase", - E2EEAlgorithm: "v2", - usePathObfuscation: true, - encryptInternalMetadata: true, - })), - }, - onConfirmApplySettingsFromWizard, + onlyE2EEConfiguration: vi.fn(async () => { + host.core.settings.encryptInternalMetadata = true; + return true; + }), }; const host = { editingSettings: { ...originalSettings }, initialSettings: { ...originalSettings }, - services: { setting: { applyPartial } }, core: { settings: { ...originalSettings }, getModule: vi.fn(() => setupManager), @@ -198,8 +189,7 @@ describe("paneRemoteConfig", () => { paneRemoteConfig.call(host as never, {} as HTMLElement, { addPanel } as never); await runtime.clickHandlers[0](); - expect(applyPartial).toHaveBeenCalledWith({ encryptInternalMetadata: true }, true); - expect(onConfirmApplySettingsFromWizard).not.toHaveBeenCalled(); + expect(setupManager.onlyE2EEConfiguration).toHaveBeenCalledOnce(); expect(host.editingSettings.encryptInternalMetadata).toBe(true); expect(host.initialSettings.encryptInternalMetadata).toBe(true); expect(host.requestUpdate).toHaveBeenCalledOnce(); diff --git a/src/modules/features/SetupManager.ts b/src/modules/features/SetupManager.ts index acc5e87a..7c6aa349 100644 --- a/src/modules/features/SetupManager.ts +++ b/src/modules/features/SetupManager.ts @@ -336,6 +336,31 @@ export class SetupManager extends AbstractModule { this._log("E2EE configuration cancelled.", LOG_LEVEL_NOTICE); 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 = { ...currentSetting, ...e2eeConf, diff --git a/src/modules/features/SetupManager.unit.spec.ts b/src/modules/features/SetupManager.unit.spec.ts index 41d9ecaf..c55c0354 100644 --- a/src/modules/features/SetupManager.unit.spec.ts +++ b/src/modules/features/SetupManager.unit.spec.ts @@ -659,3 +659,23 @@ describe("SetupManager", () => { 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(); + } + ); +}); diff --git a/src/modules/features/SetupWizard/dialogs/SetupRemoteE2EE.svelte b/src/modules/features/SetupWizard/dialogs/SetupRemoteE2EE.svelte index 1b02f79c..43ecdb2a 100644 --- a/src/modules/features/SetupWizard/dialogs/SetupRemoteE2EE.svelte +++ b/src/modules/features/SetupWizard/dialogs/SetupRemoteE2EE.svelte @@ -105,7 +105,8 @@ (Obfuscate Properties). The remote type is selected later in this setup wizard.
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. diff --git a/src/serviceFeatures/replication/ReplicateResultProcessor.ts b/src/serviceFeatures/replication/ReplicateResultProcessor.ts index 8cd69b3a..f0010c96 100644 --- a/src/serviceFeatures/replication/ReplicateResultProcessor.ts +++ b/src/serviceFeatures/replication/ReplicateResultProcessor.ts @@ -1,7 +1,7 @@ +import { assessRemoteFeatureDocument, describeRemoteFeatureRejection } from "@vrtmrz/livesync-commonlib/replication"; import { SYNCINFO_ID, VERSIONING_DOCID, - type EntryVersionInfo, type AnyEntry, type EntryDoc, type EntryLeaf, @@ -9,11 +9,6 @@ import { type MetaEntry, type ObsidianLiveSyncSettings, } 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 { LOG_LEVEL_DEBUG, @@ -69,10 +64,6 @@ interface ReplicateResultProcessorContext { readonly services: ReplicateResultProcessorServices; } type ReplicateResultProcessorState = { - databaseId?: string; - observedFeatures?: string[]; - highestObservedVersion?: number; - invalidControlObserved?: boolean; queued: PouchDB.Core.ExistingDocument[]; processing: PouchDB.Core.ExistingDocument[]; }; @@ -83,10 +74,6 @@ function shortenRev(rev: string | undefined): string { if (!rev) return "undefined"; return rev.length > 10 ? rev.substring(0, 10) : rev; } -function getPhysicalDatabaseId(database: PouchDB.Database): Promise { - const identified = database as PouchDB.Database & { id?: () => Promise }; - return typeof identified.id === "function" ? identified.id() : Promise.resolve(undefined); -} export class ReplicateResultProcessor { private log(message: string, level: LOG_LEVEL = LOG_LEVEL_INFO) { Logger(`[ReplicateResultProcessor] ${message}`, level); @@ -129,89 +116,6 @@ export class ReplicateResultProcessor { // If true, the processing queue processor bails the loop. 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; - private _observedFeatures = new Set(); - 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) { - 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. * @@ -232,8 +136,6 @@ export class ReplicateResultProcessor { public get isSuspended() { return ( this._suspended || - this._compatibilityBlocked || - this._assessingDatabase || !this.acceptsResultApplication || this.context.currentSettings().suspendParseReplicationResult || this.services.appLifecycle.isSuspended() @@ -244,38 +146,17 @@ export class ReplicateResultProcessor { * Take a snapshot of the current processing state. * This snapshot is stored in the KV database for recovery on restart. */ - private _snapshotWriter: Promise = Promise.resolve(); - - protected _takeSnapshot(): Promise { - // A blocked-batch flush must follow any earlier throttled write, or an - // older snapshot could replace the queue after the replication callback. - const write = this._snapshotWriter - .catch((): void => undefined) - .then(async () => { - const physicalDatabase = this.localDatabase.localDatabase; - const databaseId = await getPhysicalDatabaseId(physicalDatabase); - if (physicalDatabase !== this.localDatabase.localDatabase) return; - 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 { - if (this._compatibilityBlocked) await this._takeSnapshot(); + protected async _takeSnapshot() { + const snapshot = { + 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(); } /** * Trigger taking a snapshot. @@ -292,42 +173,10 @@ export class ReplicateResultProcessor { * Restore from snapshot. */ 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 .getKeyValueDB() .get(KV_KEY_REPLICATION_RESULT_PROCESSOR_SNAPSHOT); - if (physicalDatabase !== this.localDatabase.localDatabase) return; - 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)) { + if (snapshot) { // Restoring the snapshot re-runs processing for both queued and processing items. const newQueue = [...snapshot.processing, ...snapshot.queued, ...this._queuedChanges]; this._queuedChanges = []; @@ -338,9 +187,6 @@ export class ReplicateResultProcessor { ); // await this._takeSnapshot(); } - this._assessingDatabase = false; - this.updateProcessingActivity(); - this.triggerProcessQueue(); } private _restoreFromSnapshot: Promise | undefined = undefined; @@ -350,9 +196,7 @@ export class ReplicateResultProcessor { * @returns Promise that resolves when restoration is complete. */ public restoreFromSnapshotOnce() { - this.refreshPhysicalDatabase(); if (!this._restoreFromSnapshot) { - this._assessingDatabase = true; this._restoreFromSnapshot = this.restoreFromSnapshot(); } return this._restoreFromSnapshot; @@ -386,17 +230,7 @@ export class ReplicateResultProcessor { * @param changes Changes to enqueue */ - public enqueueAll(changes: PouchDB.Core.ExistingDocument[], sourceDatabase?: PouchDB.Database) { - 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); - } + public enqueueAll(changes: PouchDB.Core.ExistingDocument[]) { for (const change of changes) { // Check if the change is not a document change (e.g., chunk, versioninfo, syncinfo), and processed it directly. const isProcessed = this.processIfNonDocumentChange(change); @@ -441,8 +275,17 @@ export class ReplicateResultProcessor { this.log(`Processed chunk: ${shortenId(change._id)}`, LOG_LEVEL_DEBUG); 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); + 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; } if ( @@ -568,12 +411,11 @@ export class ReplicateResultProcessor { // (per-document serialisation caps concurrency). const releaser = await this._semaphore.acquire(); releaser(); - if (this.isSuspended) break; // Dequeue the next change const doc = this._queuedChanges.shift(); if (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) this.triggerTakeSnapshot(); @@ -589,12 +431,8 @@ export class ReplicateResultProcessor { * @param change * @returns */ - async parseDocumentChange( - change: PouchDB.Core.ExistingDocument, - sourceDatabase: PouchDB.Database = this.localDatabase.localDatabase - ) { + async parseDocumentChange(change: PouchDB.Core.ExistingDocument) { try { - if (this.shouldStopApplication(sourceDatabase)) return; if (isAnyNote(change)) { const docMtime = change.mtime ?? 0; 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 (await this.services.replication.processVirtualDocument(change)) return; - if (this.shouldStopApplication(sourceDatabase)) return; // If the document is version info, check compatibility and return. if (isAnyNote(change)) { const docPath = this.getPath(change); @@ -619,7 +456,6 @@ export class ReplicateResultProcessor { this.log(`Skipped: ${docPath}`, LOG_LEVEL_VERBOSE); return; } - if (this.shouldStopApplication(sourceDatabase)) return; const size = change.size; // Note that this size check depends size that in metadata, not the actual content size. if (this.services.vault.isFileSizeTooLarge(size)) { @@ -629,20 +465,11 @@ export class ReplicateResultProcessor { ); return; } - return await this.applyToDatabase(change, sourceDatabase); + return await this.applyToDatabase(change); } this.log(`Skipped unexpected non-note document: ${change._id}`, LOG_LEVEL_INFO); return; } 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 this._processingChanges = this._processingChanges.filter((e) => e !== change); try { @@ -662,16 +489,12 @@ export class ReplicateResultProcessor { } // Phase 2: apply the document to database - protected applyToDatabase( - doc: PouchDB.Core.ExistingDocument, - sourceDatabase: PouchDB.Database = this.localDatabase.localDatabase - ) { + protected applyToDatabase(doc: PouchDB.Core.ExistingDocument) { return this.withCounting(async () => { let releaser: Awaited> | undefined = undefined; try { releaser = await this._semaphore.acquire(); - if (this.shouldStopApplication(sourceDatabase)) return; - await this._applyToDatabase(doc, sourceDatabase); + await this._applyToDatabase(doc); } catch (e) { this.log(`Error while processing replication result`, LOG_LEVEL_NOTICE); this.logError(e); @@ -685,16 +508,12 @@ export class ReplicateResultProcessor { } // Phase 2.1: process the document and apply to storage // This function is serialized per document to avoid race-condition for the same document. - private _applyToDatabase( - doc_: PouchDB.Core.ExistingDocument, - sourceDatabase: PouchDB.Database - ) { + private _applyToDatabase(doc_: PouchDB.Core.ExistingDocument) { const dbDoc = doc_ as LoadedEntry; // It has no `data` const path = this.getPath(dbDoc); return serialized(`replication-process:${dbDoc._id}`, async () => { const docNote = `${path} (${shortenId(dbDoc._id)}, ${shortenRev(dbDoc._rev)})`; const isRequired = await this.checkIsChangeRequiredForDatabaseProcessing(dbDoc); - if (this.shouldStopApplication(sourceDatabase)) return; if (!isRequired) { this.log(`Skipped (Not latest): ${docNote}`, LOG_LEVEL_VERBOSE); return; @@ -708,7 +527,6 @@ export class ReplicateResultProcessor { const doc = isDeleted ? { ...dbDoc, data: "" } : await this.localDatabase.getDBEntryFromMeta({ ...dbDoc }, false, true); - if (this.shouldStopApplication(sourceDatabase)) return; if (!doc) { // Failed to gather content this.log(`Failed to gather content of ${docNote}`, LOG_LEVEL_NOTICE); @@ -719,10 +537,9 @@ export class ReplicateResultProcessor { // Already processed this.log(`Processed by other processor: ${docNote}`, LOG_LEVEL_DEBUG); } else if (this.services.vault.isValidPath(this.getPath(doc))) { - if (this.shouldStopApplication(sourceDatabase)) return; // Apply to storage if the path is valid try { - const reflected = await this.applyToStorage(doc as MetaEntry, sourceDatabase); + const reflected = await this.applyToStorage(doc as MetaEntry); if (!reflected) { this.reportVaultReflectionFailure(doc as MetaEntry); return; @@ -743,15 +560,9 @@ export class ReplicateResultProcessor { * @param entry * @returns */ - protected applyToStorage( - entry: MetaEntry, - sourceDatabase: PouchDB.Database = this.localDatabase.localDatabase - ) { + protected applyToStorage(entry: MetaEntry) { return this.withCounting( - () => - this.shouldStopApplication(sourceDatabase) - ? Promise.resolve(false) - : this.services.replication.processSynchroniseResult(entry), + () => this.services.replication.processSynchroniseResult(entry), this.services.replication.storageApplyingCount ); } diff --git a/src/serviceFeatures/replication/ReplicateResultProcessor.unit.spec.ts b/src/serviceFeatures/replication/ReplicateResultProcessor.unit.spec.ts index 469c4b3f..2ea388ed 100644 --- a/src/serviceFeatures/replication/ReplicateResultProcessor.unit.spec.ts +++ b/src/serviceFeatures/replication/ReplicateResultProcessor.unit.spec.ts @@ -120,6 +120,34 @@ function setup(options: SetupOptions = {}) { } 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([ ["Windows", isValidFilenameInWidows], ["Android", isValidFilenameInAndroid], @@ -249,30 +277,6 @@ describe("ReplicateResultProcessor", () => { 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; - 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 () => { const { onCloseActiveReplication, processor, processSynchroniseResult } = setup(); const versionInfo = { @@ -289,149 +293,23 @@ describe("ReplicateResultProcessor", () => { expect(onCloseActiveReplication).not.toHaveBeenCalled(); }); - it("rechecks shared writer settings when a supported feature is added to an active database", async () => { - const { onCloseActiveReplication, processor } = setup({ - localVersionInfo: { - _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; - - 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(); - 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, - ]); - 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, - 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); + it("reports unknown feature identifiers and requests Replicator retirement", () => { + const logger = vi.fn(); + setGlobalLogFunction(logger); try { - const { onCloseActiveReplication, processor } = setup(); - const versionInfo = { - _id: VERSIONING_DOCID, - _rev: "2-unknown", - type: "versioninfo", - version: REMOTE_FEATURE_GENERATION, - used_features: ["future-format-v7"], - } as PouchDB.Core.ExistingDocument; - - processor.enqueueAll([versionInfo]); - processor.enqueueAll([versionInfo]); - + const { processor, onCloseActiveReplication } = setup(); + processor.enqueueAll([ + { + _id: VERSIONING_DOCID, + _rev: "1-unknown", + type: "versioninfo", + version: REMOTE_FEATURE_GENERATION, + used_features: ["future-format-v7"], + } as PouchDB.Core.ExistingDocument, + ]); expect(onCloseActiveReplication).toHaveBeenCalledOnce(); - expect(log).toHaveBeenCalledWith( - "[ReplicateResultProcessor] Unknown features are in use: future-format-v7", + expect(logger).toHaveBeenCalledWith( + expect.stringContaining("future-format-v7"), LOG_LEVEL_NOTICE, 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; - const versionInfo = { - _id: VERSIONING_DOCID, - _rev: "2-unknown", - type: "versioninfo", - version: REMOTE_FEATURE_GENERATION, - used_features: ["future-format-v7"], - } as PouchDB.Core.ExistingDocument; - - 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, - ]); - expect(processor.isSuspended).toBe(true); - - localDatabase.localDatabase = {} as PouchDB.Database; - await processor.restoreFromSnapshotOnce(); - - expect(processor.isSuspended).toBe(false); - expect(onCloseActiveReplication).toHaveBeenCalledOnce(); - }); - it("scans normal-file metadata without loading chunk documents and requeues it", async () => { const documents = [ { _id: "first", _rev: "1-a", type: "plain", path: "first.md" }, diff --git a/src/serviceFeatures/replication/index.ts b/src/serviceFeatures/replication/index.ts index d5075d04..74a4f1e4 100644 --- a/src/serviceFeatures/replication/index.ts +++ b/src/serviceFeatures/replication/index.ts @@ -98,21 +98,19 @@ export function useReplicationFeature { - await resultProcessor.restoreFromSnapshotOnce(); - return true; + services.databaseEvents.onDatabaseInitialised.addHandler(() => { + fireAndForget(() => resultProcessor.restoreFromSnapshotOnce()); + return Promise.resolve(true); }); services.appLifecycle.onSettingLoaded.addHandler(initialiseAutomaticReplicationTriggers); - services.replication.parseSynchroniseResult.addHandler(async (documents, sourceDatabase) => { - resultProcessor.enqueueAll(documents, sourceDatabase); - await resultProcessor.persistBlockedSnapshot(); - return true; + services.replication.parseSynchroniseResult.addHandler((documents) => { + resultProcessor.enqueueAll(documents); + return Promise.resolve(true); }); services.replication.onBeforeReplicate.addHandler(onlinePreflight, 10); services.replication.onPrepareCentralRemoteReplication.addHandler(securitySeedPreflight); services.replication.onBeforeReplicate.addHandler(async () => { await resultProcessor.restoreFromSnapshotOnce(); - if (resultProcessor.isCompatibilityBlocked) return false; unresolvedErrorManager.clearErrors(); return true; }, 100); diff --git a/src/serviceFeatures/replication/replicationFeature.unit.spec.ts b/src/serviceFeatures/replication/replicationFeature.unit.spec.ts index c151b6a9..fc87dc7a 100644 --- a/src/serviceFeatures/replication/replicationFeature.unit.spec.ts +++ b/src/serviceFeatures/replication/replicationFeature.unit.spec.ts @@ -22,7 +22,9 @@ type SetupOptions = { function setup(options: SetupOptions = {}) { const defaultLocalDatabase = { localDatabase: {}, - getRaw: vi.fn(async () => { throw { status: 404 }; }), + getRaw: vi.fn(async () => { + throw { status: 404 }; + }), }; const { getLocalDatabase = () => defaultLocalDatabase, @@ -183,59 +185,4 @@ describe("replication serviceFeature composition", () => { retirement.resolve(true); }); - - it("persists a blocked batch before its replication callback settles", async () => { - const writeFinished = promiseWithResolvers(); - 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; - const unknownVersion = { - _id: VERSIONING_DOCID, - _rev: "2-unknown", - type: "versioninfo", - version: REMOTE_FEATURE_GENERATION, - used_features: ["future-format-v7"], - } as PouchDB.Core.ExistingDocument; - 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; - - await parseHandler!([unknownVersion]); - - await expect(beforeReplicateHandlers.get(100)!(false)).resolves.toBe(false); - }); }); diff --git a/test/e2e-obsidian/scripts/hidden-file-snippet-sync.ts b/test/e2e-obsidian/scripts/hidden-file-snippet-sync.ts index fe48c6fe..5079c021 100644 --- a/test/e2e-obsidian/scripts/hidden-file-snippet-sync.ts +++ b/test/e2e-obsidian/scripts/hidden-file-snippet-sync.ts @@ -368,13 +368,20 @@ async function uploadHiddenFile( return ids.has(entry.id) && entry.children.every((childId) => ids.has(childId)); }); const remoteEntry = await fetchCouchDbDocument(context.couchDb, context.dbName, entry.id); - if (!remoteEntry.path?.startsWith("/\\:") || remoteEntry.children?.length !== 0 || - remoteEntry.ctime !== 0 || remoteEntry.mtime !== 0 || remoteEntry.size !== 0) { + if ( + !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}.`); } const versionInfo = await fetchCouchDbDocument(context.couchDb, context.dbName, VERSIONING_DOCID); - if (versionInfo.version !== REMOTE_FEATURE_GENERATION || - !(versionInfo.used_features as unknown[] | undefined)?.includes(ENCRYPTED_INTERNAL_METADATA_FEATURE)) { + if ( + 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."); } return entry; @@ -400,6 +407,29 @@ async function runCreateRoundTrip( await writeVaultFile(vaultA.path, snippetPath, snippetContent); let session = await startConfiguredSession(context, vaultA); 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(); session = await startConfiguredSession(context, vaultB); @@ -588,11 +618,14 @@ async function runInitialisationNoticeGrouping(context: RunnerContext, vault: Te await withObsidianPage(port, async (page) => { const deadline = Date.now() + timeoutMs; while ((await page.locator(".notice:visible").count()) > 0 && Date.now() < deadline) { - await page.locator(".notice:visible").first().click({ - force: true, - position: { x: 2, y: 2 }, - timeout: timeoutMs, - }); + await page + .locator(".notice:visible") + .first() + .click({ + force: true, + position: { x: 2, y: 2 }, + timeout: timeoutMs, + }); } assertEqual( await page.locator(".notice:visible").count(), @@ -724,17 +757,15 @@ async function runInitialisationNoticeGrouping(context: RunnerContext, vault: Te const result = await withObsidianPage(port, async (page) => { await page.evaluate((stateKey) => { - const state = (globalThis as unknown as Record< - string, - { releasePreparation?: () => void } | undefined - >)[stateKey]; + const state = ( + globalThis as unknown as Record void } | undefined> + )[stateKey]; state?.releasePreparation?.(); }, hiddenFileInitialisationStateKey); await page.waitForFunction( (stateKey) => - (globalThis as unknown as Record)[ - stateKey - ]?.reachedInitialisation === true, + (globalThis as unknown as Record)[stateKey] + ?.reachedInitialisation === true, hiddenFileInitialisationStateKey, { timeout: timeoutMs } ); diff --git a/test/e2e-obsidian/scripts/remote-feature-change.ts b/test/e2e-obsidian/scripts/remote-feature-change.ts index afdf945a..8e9b66b3 100644 --- a/test/e2e-obsidian/scripts/remote-feature-change.ts +++ b/test/e2e-obsidian/scripts/remote-feature-change.ts @@ -32,9 +32,6 @@ const unknownFeature = "future-format-v7"; type FeatureState = { version: number | null; features: string[]; - snapshotFeatures: string[]; - snapshotDatabaseId: string | null; - databaseId: string | null; hasActiveReplicator: boolean; }; @@ -46,14 +43,9 @@ async function readFeatureState(cliBinary: string, env: NodeJS.ProcessEnv): Prom "const core=app.plugins.plugins['obsidian-livesync'].core;", `const id=${JSON.stringify(VERSIONING_DOCID)};`, "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({", "version:typeof info?.version==='number'?info.version:null,", "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(),", "});", "})()", @@ -153,12 +145,8 @@ async function main(): Promise { const observed = await waitForState( cli.binary, session.cliEnv, - (state) => - state.version === 13 && - state.features.includes(unknownFeature) && - state.snapshotFeatures.includes(unknownFeature) && - !state.hasActiveReplicator, - "the live feature change, durable observation, and Replicator retirement" + (state) => state.version === 13 && state.features.includes(unknownFeature) && !state.hasActiveReplicator, + "the live feature change and Replicator retirement" ); const replicated = await evalObsidianJson( @@ -171,25 +159,6 @@ async function main(): Promise { if (acceptedAfterStop !== acceptedContent) 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( - 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(); session = undefined; session = await startObsidianLiveSyncSession({ binary, cliBinary: cli.binary, vault }); @@ -197,36 +166,40 @@ async function main(): Promise { const afterRestart = await waitForState( cli.binary, session.cliEnv, - (state) => - state.version === 13 && state.features.length === 0 && state.snapshotFeatures.includes(unknownFeature), - "the retained unknown feature after restart" + (state) => state.version === 13 && state.features.includes(unknownFeature), + "the remote feature requirement after restart" ); - if (afterRestart.snapshotDatabaseId !== afterRestart.databaseId) - throw new Error("The retained feature observation belongs to a different physical database."); const replicatedAfterRestart = await evalObsidianJson( cli.binary, "(async()=>JSON.stringify(!!(await app.plugins.plugins['obsidian-livesync'].core.services.replication.replicate(true))))()", session.cliEnv ); - const continuousAfterRestart = await evalObsidianJson<{ status: string }>( + const continuousAfterRestart = await evalObsidianJson<{ requestStatus: string; connected: boolean }>( cli.binary, [ "(async()=>{", "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'}});", - "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(""), session.cliEnv ); - if (replicatedAfterRestart || continuousAfterRestart.status !== "blocked") + if (replicatedAfterRestart || continuousAfterRestart.connected !== false) 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) throw new Error("Previously accepted Vault content changed after restart."); 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 { await session?.app.stop(); diff --git a/updates.md b/updates.md index 7be7477a..151f8806 100644 --- a/updates.md +++ b/updates.md @@ -17,7 +17,7 @@ Earlier releases remain available in the 1.0 release history, the 1.0 preview hi #### 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. -- 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