import { promiseWithResolvers } from "octagonal-wheels/promises"; import { reactiveSource } from "octagonal-wheels/dataobject/reactive"; import { describe, expect, it, vi } from "vitest"; import { VERSIONING_DOCID, type EntryDoc, type FilePathWithPrefix, } from "@vrtmrz/livesync-commonlib/compat/common/types"; import { ENCRYPTED_INTERNAL_METADATA_FEATURE, REMOTE_FEATURE_GENERATION } from "@vrtmrz/livesync-commonlib/replication"; import { isValidFilenameInAndroid, isValidFilenameInWidows, } from "@vrtmrz/livesync-commonlib/compat/string_and_binary/path"; import { defaultLogger, LOG_LEVEL_DEBUG, LOG_LEVEL_NOTICE, LOG_LEVEL_VERBOSE, setGlobalLogFunction, } from "octagonal-wheels/common/logger"; import { ReplicateResultProcessor } from "./ReplicateResultProcessor"; function note(id: string): PouchDB.Core.ExistingDocument { return { _id: id, _rev: "1-test", path: `${id}.md`, ctime: 1, mtime: 2, size: 1, children: [], datatype: "plain", type: "plain", eden: {}, } as unknown as PouchDB.Core.ExistingDocument; } type SetupOptions = { applicationReady?: boolean; databaseReady?: boolean; maxMTimeForReflectEvents?: number; isValidPath?: (path: string) => boolean; processSynchroniseResult?: (entry: unknown) => Promise; setSnapshot?: (key: string, value: unknown) => Promise; getSnapshot?: (key: string) => Promise; localVersionInfo?: unknown; isTargetFile?: (path: string) => Promise; databaseId?: string; }; function setup(options: SetupOptions = {}) { const processSynchroniseResult = vi.fn(options.processSynchroniseResult ?? (async () => true)); const setSnapshot = vi.fn(options.setSnapshot ?? (async () => undefined)); const runBoundedLocalApplicationActivity = vi.fn(async (task: () => Promise) => await task()); const onCloseActiveReplication = vi.fn(async () => true); const isReady = vi.fn(() => options.applicationReady ?? true); const isValidPath = vi.fn(options.isValidPath ?? (() => true)); const getDBEntryFromMeta = vi.fn(async (entry: object) => ({ ...entry, data: "x" })); const localPhysicalDatabase = { ...(options.databaseId ? { id: vi.fn(async () => options.databaseId) } : {}), } as PouchDB.Database; const core = { services: { appLifecycle: { isReady, isSuspended: () => false }, database: { isDatabaseReady: () => options.databaseReady ?? true }, path: { getPath: (entry: { path: string }) => entry.path }, replication: { databaseQueueCount: reactiveSource(0), storageApplyingCount: reactiveSource(0), replicationResultCount: reactiveSource(0), processVirtualDocument: vi.fn(async () => false), processOptionalSynchroniseResult: vi.fn(async () => false), processSynchroniseResult, }, replicator: { onCloseActiveReplication, runBoundedLocalApplicationActivity }, vault: { isTargetFile: vi.fn(options.isTargetFile ?? (async () => true)), isFileSizeTooLarge: vi.fn(() => false), isValidPath, }, }, kvDB: { set: setSnapshot, get: vi.fn(options.getSnapshot ?? (async () => undefined)) }, localDatabase: { localDatabase: localPhysicalDatabase, getRaw: vi.fn(async (id: string) => { if (id === VERSIONING_DOCID) { if (options.localVersionInfo === undefined) throw { status: 404 }; return options.localVersionInfo; } return { _id: id, _rev: "1-test" }; }), getDBEntryFromMeta, }, }; const processor = new ReplicateResultProcessor({ currentSettings: () => ({ maxMTimeForReflectEvents: options.maxMTimeForReflectEvents ?? 0, suspendParseReplicationResult: false, }), getKeyValueDB: () => core.kvDB, getLocalDatabase: () => core.localDatabase, requestActiveReplicatorRetirement: () => { void onCloseActiveReplication(); }, runLocalApplicationActivity: runBoundedLocalApplicationActivity, services: core.services, } as never); return { getDBEntryFromMeta, isTargetFile: core.services.vault.isTargetFile, localPhysicalDatabase, localDatabase: core.localDatabase, isReady, isValidPath, onCloseActiveReplication, processor, processSynchroniseResult, runBoundedLocalApplicationActivity, }; } 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], ])("does not reflect a replicated colon path into the %s Vault", async (_platform, validatePath) => { const path = "Folder/Poem: Example.md" as FilePathWithPrefix; const document = { ...note("colon-path"), path }; const { getDBEntryFromMeta, isValidPath, processor, processSynchroniseResult } = setup({ isValidPath: validatePath, }); processor.enqueueAll([document]); await vi.waitFor(() => expect(isValidPath).toHaveBeenCalledWith(path)); await vi.waitFor(() => expect(processor["_processingChanges"]).toHaveLength(0)); expect(getDBEntryFromMeta).toHaveBeenCalledWith(expect.objectContaining({ path }), false, true); expect(processSynchroniseResult).not.toHaveBeenCalled(); }); it("resumes another document after in-flight updates to one document fill the application slots", async () => { const hotGate = promiseWithResolvers(); const { processor, processSynchroniseResult } = setup({ processSynchroniseResult: async (entry) => { if ((entry as { _id: string })._id === "hot-queue") return await hotGate.promise; return true; }, }); try { for (let index = 1; index <= 10; index++) { // A queued duplicate is coalesced; a new notification for a document // already being processed can occupy another application slot. processor.enqueueAll([note("hot-queue")]); await vi.waitFor(() => expect(processor["_processingChanges"]).toHaveLength(index)); } processor.enqueueAll([note("unrelated-queue")]); await vi.waitFor(() => { expect(processor["_semaphore"].waiting).toBeGreaterThan(0); expect(processSynchroniseResult).toHaveBeenCalledTimes(1); }); expect(processor["_queuedChanges"].map((entry) => entry._id)).toEqual(["unrelated-queue"]); } finally { hotGate.resolve(true); await vi.waitFor(() => { expect(processor["_processingChanges"]).toHaveLength(0); expect(processor["_queuedChanges"]).toHaveLength(0); }); } expect(processSynchroniseResult).toHaveBeenCalledTimes(11); expect( processSynchroniseResult.mock.calls.some(([entry]) => (entry as { _id: string })._id === "unrelated-queue") ).toBe(true); }); it("suspends result application while the application is not ready", () => { const { isReady, processor } = setup({ applicationReady: false }); expect(processor.isSuspended).toBe(true); expect(isReady).toHaveBeenCalledOnce(); }); it("applies documents held before readiness once the application becomes ready", async () => { const { isReady, processor, processSynchroniseResult } = setup({ applicationReady: false }); processor.enqueueAll([note("held")]); await new Promise((resolve) => setTimeout(resolve, 0)); expect(processSynchroniseResult).not.toHaveBeenCalled(); expect(processor["_queuedChanges"]).toHaveLength(1); isReady.mockReturnValue(true); processor.continueHeldDocuments(); await vi.waitFor(() => expect(processSynchroniseResult).toHaveBeenCalledOnce()); expect(processor["_queuedChanges"]).toHaveLength(0); }); it("keeps an explicit suspension when the application becomes ready", async () => { const { isReady, processor, processSynchroniseResult } = setup({ applicationReady: false }); processor.suspend(); processor.enqueueAll([note("held")]); isReady.mockReturnValue(true); processor.continueHeldDocuments(); await new Promise((resolve) => setTimeout(resolve, 0)); expect(processSynchroniseResult).not.toHaveBeenCalled(); expect(processor["_queuedChanges"]).toHaveLength(1); processor.resume(); await vi.waitFor(() => expect(processSynchroniseResult).toHaveBeenCalledOnce()); }); it("applies results in remediation mode, which never reports readiness", () => { const { processor } = setup({ applicationReady: false, maxMTimeForReflectEvents: Date.parse("2026-09-01T00:00:00Z"), }); expect(processor.isSuspended).toBe(false); }); it("holds results in remediation mode while the local database is being rebuilt", () => { const { processor } = setup({ applicationReady: false, databaseReady: false, maxMTimeForReflectEvents: Date.parse("2026-09-01T00:00:00Z"), }); expect(processor.isSuspended).toBe(true); }); it("still skips a document modified after the limit while the application is unready", async () => { const maxMTimeForReflectEvents = Date.parse("2026-09-01T00:00:00Z"); const { processor, processSynchroniseResult } = setup({ applicationReady: false, maxMTimeForReflectEvents, }); const tooRecent = { ...note("too-recent"), mtime: maxMTimeForReflectEvents + 1, } as PouchDB.Core.ExistingDocument; processor.enqueueAll([tooRecent]); await vi.waitFor(() => { expect(processor["_queuedChanges"]).toHaveLength(0); expect(processor["_processingChanges"]).toHaveLength(0); }); expect(processSynchroniseResult).not.toHaveBeenCalled(); }); it("retires active ownership when a newer remote version is observed", async () => { const { onCloseActiveReplication, processor } = setup(); const versionInfo = { _id: VERSIONING_DOCID, _rev: "1-test", type: "versioninfo", version: REMOTE_FEATURE_GENERATION + 1, } as unknown as PouchDB.Core.ExistingDocument; processor.enqueueAll([versionInfo]); await vi.waitFor(() => expect(onCloseActiveReplication).toHaveBeenCalledOnce()); }); it("continues applying documents after restoring a legacy local version document", async () => { const { onCloseActiveReplication, processor, processSynchroniseResult } = setup({ localVersionInfo: { _id: VERSIONING_DOCID, type: "versioninfo", version: 11, }, }); await processor.restoreFromSnapshotOnce(); processor.enqueueAll([note("legacy-database-note")]); await vi.waitFor(() => expect(processSynchroniseResult).toHaveBeenCalledOnce()); expect(processor.isSuspended).toBe(false); expect(onCloseActiveReplication).not.toHaveBeenCalled(); }); it("continues when a newly received feature is supported", async () => { const { onCloseActiveReplication, processor, processSynchroniseResult } = setup(); 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, note("supported-update")]); await vi.waitFor(() => expect(processSynchroniseResult).toHaveBeenCalledOnce()); expect(onCloseActiveReplication).not.toHaveBeenCalled(); }); it("reports unknown feature identifiers and requests Replicator retirement", () => { const logger = vi.fn(); setGlobalLogFunction(logger); try { 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(logger).toHaveBeenCalledWith( expect.stringContaining("future-format-v7"), LOG_LEVEL_NOTICE, undefined ); } finally { setGlobalLogFunction(defaultLogger); } }); 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" }, { _id: "second", _rev: "1-b", type: "plain", path: "second.md" }, ] as unknown as PouchDB.Core.ExistingDocument[]; const findAllNormalDocs = vi.fn(async function* () { yield* documents; }); const getLocalDatabase = vi.fn(() => ({ findAllNormalDocs })); const processor = new ReplicateResultProcessor({ getLocalDatabase, } as never); const enqueueAll = vi.spyOn(processor, "enqueueAll").mockImplementation(() => undefined); await expect(processor.reprocessStoredDocuments()).resolves.toBe(2); expect(findAllNormalDocs).toHaveBeenCalledOnce(); expect(getLocalDatabase).toHaveBeenCalledOnce(); expect(enqueueAll).toHaveBeenCalledOnce(); expect(enqueueAll).toHaveBeenCalledWith(documents); }); it("keeps one local application activity until every replicated document has been applied", async () => { const applying = promiseWithResolvers(); let activityFinished = false; const { processor, processSynchroniseResult, runBoundedLocalApplicationActivity } = setup({ processSynchroniseResult: async () => applying.promise, }); runBoundedLocalApplicationActivity.mockImplementation(async (task: () => Promise) => { await task(); activityFinished = true; }); processor.enqueueAll([note("one"), note("two")]); await vi.waitFor(() => expect(processSynchroniseResult).toHaveBeenCalledTimes(2)); expect(runBoundedLocalApplicationActivity).toHaveBeenCalledTimes(1); expect(runBoundedLocalApplicationActivity).toHaveBeenCalledWith(expect.any(Function), { label: "replicated-document-application", }); expect(activityFinished).toBe(false); applying.resolve(true); await vi.waitFor(() => expect(activityFinished).toBe(true)); }); it("settles local application activity when the final recovery snapshot fails", async () => { let activityFinished = false; const { processor, runBoundedLocalApplicationActivity } = setup({ setSnapshot: async () => Promise.reject(new Error("snapshot failed")), }); runBoundedLocalApplicationActivity.mockImplementation(async (task: () => Promise) => { await task(); activityFinished = true; }); processor.enqueueAll([note("one")]); await vi.waitFor(() => expect(activityFinished).toBe(true)); }); it("releases and reacquires local application activity around processing suspension", async () => { const applying = promiseWithResolvers(); let completedActivities = 0; const { processor, processSynchroniseResult, runBoundedLocalApplicationActivity } = setup({ processSynchroniseResult: async () => applying.promise, }); runBoundedLocalApplicationActivity.mockImplementation(async (task: () => Promise) => { await task(); completedActivities++; }); processor.enqueueAll([note("one")]); await vi.waitFor(() => expect(processSynchroniseResult).toHaveBeenCalledOnce()); processor.suspend(); await vi.waitFor(() => expect(completedActivities).toBe(1)); processor.resume(); await vi.waitFor(() => expect(runBoundedLocalApplicationActivity).toHaveBeenCalledTimes(2)); applying.resolve(true); await vi.waitFor(() => expect(completedActivities).toBe(2)); }); it.each([ ["returns false", async () => false, undefined], ["throws", async () => Promise.reject(new Error("File name too long")), "File name too long"], ])("reports when Vault reflection %s", async (_description, processSynchroniseResult, errorMessage) => { const log = vi.fn((_message: unknown, _level?: number) => undefined); setGlobalLogFunction(log); try { const { processor } = setup({ processSynchroniseResult }); processor.enqueueAll([note("unreflectable")]); await vi.waitFor(() => expect(log).toHaveBeenCalledWith( "Not all files could be synchronised. Check the affected files. Generate a report to review the detailed log.", LOG_LEVEL_NOTICE, undefined ) ); expect(log).toHaveBeenCalledWith( "[ReplicateResultProcessor] Live replication could not reflect unreflectable.md from the local database to the Vault; this path remains eligible for a later Vault scan.", LOG_LEVEL_VERBOSE, undefined ); if (errorMessage !== undefined) { expect(log).toHaveBeenCalledWith( expect.objectContaining({ message: errorMessage }), LOG_LEVEL_VERBOSE, undefined ); } expect(log).not.toHaveBeenCalledWith( expect.stringContaining("Processed: unreflectable.md"), LOG_LEVEL_DEBUG, undefined ); } finally { setGlobalLogFunction(defaultLogger); } }); });