Integrate encrypted internal metadata in LiveSync

This commit is contained in:
vorotamoroz
2026-09-27 09:51:53 +00:00
parent 7b3b6ff854
commit 6c50505096
23 changed files with 1220 additions and 61 deletions
@@ -1,6 +1,7 @@
import {
SYNCINFO_ID,
VER,
VERSIONING_DOCID,
type EntryVersionInfo,
type AnyEntry,
type EntryDoc,
type EntryLeaf,
@@ -8,6 +9,11 @@ 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,
@@ -63,6 +69,10 @@ interface ReplicateResultProcessorContext {
readonly services: ReplicateResultProcessorServices;
}
type ReplicateResultProcessorState = {
databaseId?: string;
observedFeatures?: string[];
highestObservedVersion?: number;
invalidControlObserved?: boolean;
queued: PouchDB.Core.ExistingDocument<EntryDoc>[];
processing: PouchDB.Core.ExistingDocument<EntryDoc>[];
};
@@ -73,6 +83,10 @@ function shortenRev(rev: string | undefined): string {
if (!rev) return "undefined";
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 {
private log(message: string, level: LOG_LEVEL = LOG_LEVEL_INFO) {
Logger(`[ReplicateResultProcessor] ${message}`, level);
@@ -115,6 +129,89 @@ 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<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.
*
@@ -135,6 +232,8 @@ export class ReplicateResultProcessor {
public get isSuspended() {
return (
this._suspended ||
this._compatibilityBlocked ||
this._assessingDatabase ||
!this.acceptsResultApplication ||
this.context.currentSettings().suspendParseReplicationResult ||
this.services.appLifecycle.isSuspended()
@@ -145,17 +244,38 @@ export class ReplicateResultProcessor {
* Take a snapshot of the current processing state.
* This snapshot is stored in the KV database for recovery on restart.
*/
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();
private _snapshotWriter: Promise<void> = Promise.resolve();
protected _takeSnapshot(): Promise<void> {
// 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<void> {
if (this._compatibilityBlocked) await this._takeSnapshot();
}
/**
* Trigger taking a snapshot.
@@ -172,10 +292,42 @@ 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<ReplicateResultProcessorState>(KV_KEY_REPLICATION_RESULT_PROCESSOR_SNAPSHOT);
if (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)) {
// Restoring the snapshot re-runs processing for both queued and processing items.
const newQueue = [...snapshot.processing, ...snapshot.queued, ...this._queuedChanges];
this._queuedChanges = [];
@@ -186,6 +338,9 @@ export class ReplicateResultProcessor {
);
// await this._takeSnapshot();
}
this._assessingDatabase = false;
this.updateProcessingActivity();
this.triggerProcessQueue();
}
private _restoreFromSnapshot: Promise<void> | undefined = undefined;
@@ -195,7 +350,9 @@ 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;
@@ -229,7 +386,17 @@ export class ReplicateResultProcessor {
* @param changes Changes to enqueue
*/
public enqueueAll(changes: PouchDB.Core.ExistingDocument<EntryDoc>[]) {
public enqueueAll(changes: PouchDB.Core.ExistingDocument<EntryDoc>[], sourceDatabase?: PouchDB.Database<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) {
// Check if the change is not a document change (e.g., chunk, versioninfo, syncinfo), and processed it directly.
const isProcessed = this.processIfNonDocumentChange(change);
@@ -274,16 +441,8 @@ export class ReplicateResultProcessor {
this.log(`Processed chunk: ${shortenId(change._id)}`, LOG_LEVEL_DEBUG);
return true;
}
if (change.type == "versioninfo") {
if (change._id === VERSIONING_DOCID) {
this.log(`Version info document received: ${change._id}`, LOG_LEVEL_VERBOSE);
if (change.version > VER) {
// Fence and retire the active publication through its owner.
this.context.requestActiveReplicatorRetirement();
this.log(
`Remote database updated to incompatible version. update your Self-hosted LiveSync plugin.`,
LOG_LEVEL_NOTICE
);
}
return true;
}
if (
@@ -409,11 +568,12 @@ 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);
void this.parseDocumentChange(doc, this.localDatabase.localDatabase);
}
// Take snapshot (to be restored on next startup if needed)
this.triggerTakeSnapshot();
@@ -429,8 +589,12 @@ export class ReplicateResultProcessor {
* @param change
* @returns
*/
async parseDocumentChange(change: PouchDB.Core.ExistingDocument<EntryDoc>) {
async parseDocumentChange(
change: PouchDB.Core.ExistingDocument<EntryDoc>,
sourceDatabase: PouchDB.Database<EntryDoc> = this.localDatabase.localDatabase
) {
try {
if (this.shouldStopApplication(sourceDatabase)) return;
if (isAnyNote(change)) {
const docMtime = change.mtime ?? 0;
const maxMTime = this.context.currentSettings().maxMTimeForReflectEvents;
@@ -447,6 +611,7 @@ 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);
@@ -454,6 +619,7 @@ 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)) {
@@ -463,11 +629,20 @@ export class ReplicateResultProcessor {
);
return;
}
return await this.applyToDatabase(change);
return await this.applyToDatabase(change, sourceDatabase);
}
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 {
@@ -487,12 +662,16 @@ export class ReplicateResultProcessor {
}
// Phase 2: apply the document to database
protected applyToDatabase(doc: PouchDB.Core.ExistingDocument<AnyEntry>) {
protected applyToDatabase(
doc: PouchDB.Core.ExistingDocument<AnyEntry>,
sourceDatabase: PouchDB.Database<EntryDoc> = this.localDatabase.localDatabase
) {
return this.withCounting(async () => {
let releaser: Awaited<ReturnType<typeof this._semaphore.acquire>> | undefined = undefined;
try {
releaser = await this._semaphore.acquire();
await this._applyToDatabase(doc);
if (this.shouldStopApplication(sourceDatabase)) return;
await this._applyToDatabase(doc, sourceDatabase);
} catch (e) {
this.log(`Error while processing replication result`, LOG_LEVEL_NOTICE);
this.logError(e);
@@ -506,12 +685,16 @@ 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<AnyEntry>) {
private _applyToDatabase(
doc_: PouchDB.Core.ExistingDocument<AnyEntry>,
sourceDatabase: PouchDB.Database<EntryDoc>
) {
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;
@@ -525,6 +708,7 @@ 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);
@@ -535,9 +719,10 @@ 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);
const reflected = await this.applyToStorage(doc as MetaEntry, sourceDatabase);
if (!reflected) {
this.reportVaultReflectionFailure(doc as MetaEntry);
return;
@@ -558,9 +743,15 @@ export class ReplicateResultProcessor {
* @param entry
* @returns
*/
protected applyToStorage(entry: MetaEntry) {
protected applyToStorage(
entry: MetaEntry,
sourceDatabase: PouchDB.Database<EntryDoc> = this.localDatabase.localDatabase
) {
return this.withCounting(
() => this.services.replication.processSynchroniseResult(entry),
() =>
this.shouldStopApplication(sourceDatabase)
? Promise.resolve(false)
: this.services.replication.processSynchroniseResult(entry),
this.services.replication.storageApplyingCount
);
}