mirror of
https://github.com/vrtmrz/obsidian-livesync.git
synced 2026-10-08 18:32:30 +00:00
Merge current main for Setup URI dependency alignment
This commit is contained in:
@@ -10,7 +10,7 @@ import {
|
||||
import { scheduleTask } from "octagonal-wheels/concurrency/task";
|
||||
import { fireAndForget, isDirty, throttle } from "@vrtmrz/livesync-commonlib/compat/common/utils";
|
||||
import {
|
||||
collectingChunks,
|
||||
chunkFetchCounts,
|
||||
pluginScanningCount,
|
||||
hiddenFilesEventCount,
|
||||
hiddenFilesProcessingCount,
|
||||
@@ -36,7 +36,11 @@ import {
|
||||
formatRemoteActivityStatusLabel,
|
||||
getTrackedRequestCount,
|
||||
} from "./RemoteActivityStatus.ts";
|
||||
import { createMinimumVisibleActivityCount, createPaddedCounterLabel } from "./StatusBarDisplay.ts";
|
||||
import {
|
||||
createChunkFetchCounterLabel,
|
||||
createMinimumVisibleActivityCount,
|
||||
createPaddedCounterLabel,
|
||||
} from "./StatusBarDisplay.ts";
|
||||
import type { LiveSyncCore } from "@/main.ts";
|
||||
import { LiveSyncError } from "@vrtmrz/livesync-commonlib/compat/common/LSError";
|
||||
import { isValidPath } from "@/common/utils.ts";
|
||||
@@ -140,7 +144,7 @@ export class ModuleLog extends AbstractObsidianModule {
|
||||
const labelStorageCount = registerDisplay(
|
||||
createPaddedCounterLabel(this.services.replication.storageApplyingCount, `💾`)
|
||||
);
|
||||
const labelChunkCount = registerDisplay(createPaddedCounterLabel(collectingChunks, `🧩`));
|
||||
const labelChunkCount = registerDisplay(createChunkFetchCounterLabel(chunkFetchCounts));
|
||||
const labelPluginScanCount = registerDisplay(createPaddedCounterLabel(pluginScanningCount, `🔌`));
|
||||
const labelConflictProcessCount = registerDisplay(
|
||||
createPaddedCounterLabel(this.services.conflict.conflictProcessQueueCount, `🔩`)
|
||||
|
||||
@@ -131,3 +131,50 @@ export function createPaddedCounterLabel(
|
||||
source.offChanged(update);
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Displays the disjoint initial and retry chunk-fetch counts with the same
|
||||
* padding and inactive linger behaviour as the other status counters.
|
||||
*/
|
||||
export function createChunkFetchCounterLabel(
|
||||
source: ReactiveValue<{ initial: number; retrying: number }>
|
||||
): DisposableReactiveValue<string> {
|
||||
const initialCount = reactiveSource(0);
|
||||
const retryingCount = reactiveSource(0);
|
||||
const initialLabel = createPaddedCounterLabel(initialCount, "🛄");
|
||||
const retryingLabel = createPaddedCounterLabel(retryingCount, "🔁");
|
||||
const formatted = reactiveSource(`${initialLabel.value}${retryingLabel.value}`);
|
||||
let updatingCounts = false;
|
||||
let disposed = false;
|
||||
|
||||
const updateLabel = () => {
|
||||
if (updatingCounts || disposed) return;
|
||||
formatted.value = `${initialLabel.value}${retryingLabel.value}`;
|
||||
};
|
||||
initialLabel.onChanged(updateLabel);
|
||||
retryingLabel.onChanged(updateLabel);
|
||||
|
||||
const updateCounts = () => {
|
||||
if (disposed) return;
|
||||
updatingCounts = true;
|
||||
try {
|
||||
initialCount.value = source.value.initial;
|
||||
retryingCount.value = source.value.retrying;
|
||||
} finally {
|
||||
updatingCounts = false;
|
||||
updateLabel();
|
||||
}
|
||||
};
|
||||
source.onChanged(updateCounts);
|
||||
updateCounts();
|
||||
|
||||
return asDisposableReactiveValue(formatted, () => {
|
||||
if (disposed) return;
|
||||
disposed = true;
|
||||
source.offChanged(updateCounts);
|
||||
initialLabel.offChanged(updateLabel);
|
||||
retryingLabel.offChanged(updateLabel);
|
||||
initialLabel.dispose();
|
||||
retryingLabel.dispose();
|
||||
});
|
||||
}
|
||||
|
||||
@@ -3,6 +3,7 @@ import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
||||
|
||||
import {
|
||||
STATUS_COUNTER_INACTIVE_LINGER_MS,
|
||||
createChunkFetchCounterLabel,
|
||||
createMinimumVisibleActivityCount,
|
||||
createPaddedCounterLabel,
|
||||
} from "./StatusBarDisplay.ts";
|
||||
@@ -137,3 +138,41 @@ describe("createPaddedCounterLabel", () => {
|
||||
expect(display.value).toBe(" 📄\u20070");
|
||||
});
|
||||
});
|
||||
|
||||
describe("createChunkFetchCounterLabel", () => {
|
||||
beforeEach(() => {
|
||||
vi.useFakeTimers();
|
||||
});
|
||||
|
||||
afterEach(() => {
|
||||
vi.useRealTimers();
|
||||
});
|
||||
|
||||
it("keeps initial and retry counts separate across an unchanged-total handoff", () => {
|
||||
const counts = reactiveSource({ initial: 2, retrying: 5 });
|
||||
const display = createChunkFetchCounterLabel(counts);
|
||||
const withoutPadding = () => display.value.replace(/\u2007/g, "");
|
||||
const transitionSnapshots: string[] = [];
|
||||
const observeTransitions = () => transitionSnapshots.push(withoutPadding());
|
||||
|
||||
expect(withoutPadding()).toBe(" 🛄2 🔁5");
|
||||
|
||||
display.onChanged(observeTransitions);
|
||||
counts.value = { initial: 0, retrying: 7 };
|
||||
expect(withoutPadding()).toBe(" 🛄0 🔁7");
|
||||
expect(transitionSnapshots).toEqual([" 🛄0 🔁7"]);
|
||||
display.offChanged(observeTransitions);
|
||||
vi.advanceTimersByTime(STATUS_COUNTER_INACTIVE_LINGER_MS - 1);
|
||||
expect(withoutPadding()).toBe(" 🛄0 🔁7");
|
||||
vi.advanceTimersByTime(1);
|
||||
expect(withoutPadding()).toBe(" 🔁7");
|
||||
|
||||
counts.value = { initial: 0, retrying: 0 };
|
||||
expect(withoutPadding()).toBe(" 🔁0");
|
||||
display.dispose();
|
||||
vi.advanceTimersByTime(STATUS_COUNTER_INACTIVE_LINGER_MS);
|
||||
counts.value = { initial: 1, retrying: 0 };
|
||||
|
||||
expect(withoutPadding()).toBe(" 🔁0");
|
||||
});
|
||||
});
|
||||
|
||||
@@ -108,8 +108,15 @@ export class ReplicateResultProcessor {
|
||||
}
|
||||
public resume() {
|
||||
this._suspended = false;
|
||||
this.continueHeldDocuments();
|
||||
}
|
||||
/**
|
||||
* Continue the queued documents which were held, for example while the application was not ready.
|
||||
* An explicit suspension, by `suspend()` or by the settings, remains in effect.
|
||||
*/
|
||||
public continueHeldDocuments() {
|
||||
this.updateProcessingActivity();
|
||||
fireAndForget(() => this.runProcessQueue());
|
||||
this.triggerProcessQueue();
|
||||
}
|
||||
|
||||
// Whether the processing is suspended
|
||||
|
||||
@@ -207,6 +207,35 @@ describe("ReplicateResultProcessor", () => {
|
||||
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,
|
||||
|
||||
@@ -4,6 +4,7 @@ import { UnresolvedErrorManager } from "@vrtmrz/livesync-commonlib/compat/servic
|
||||
import type { ServiceContext } from "@vrtmrz/livesync-commonlib/context";
|
||||
import { fireAndForget } from "octagonal-wheels/promises";
|
||||
import type { IMinimumLiveSyncCommands, LiveSyncBaseCore } from "@/LiveSyncBaseCore";
|
||||
import { EVENT_APPLICATION_READY } from "@/common/events";
|
||||
import { createAutomaticReplicationTriggers } from "./automaticTriggers";
|
||||
import { createCentralCompatibilityRecovery } from "./centralCompatibilityRecovery";
|
||||
import { createOnlineReplicationPreflight, createSecuritySeedPreflight } from "./preflight";
|
||||
@@ -102,6 +103,9 @@ export function useReplicationFeature<TContext extends ServiceContext, TCommands
|
||||
fireAndForget(() => resultProcessor.restoreFromSnapshotOnce());
|
||||
return Promise.resolve(true);
|
||||
});
|
||||
// Commonlib emits this each time it establishes readiness. Documents held until then, such as those restored
|
||||
// from the snapshot or received during a fetch, continue from here.
|
||||
services.context.events.onEvent(EVENT_APPLICATION_READY, () => resultProcessor.continueHeldDocuments());
|
||||
services.appLifecycle.onSettingLoaded.addHandler(initialiseAutomaticReplicationTriggers);
|
||||
services.replication.parseSynchroniseResult.addHandler((documents) => {
|
||||
resultProcessor.enqueueAll(documents);
|
||||
|
||||
@@ -0,0 +1,171 @@
|
||||
import { createServiceContext } from "@vrtmrz/livesync-commonlib/context";
|
||||
import { VERSIONING_DOCID, type EntryDoc } from "@vrtmrz/livesync-commonlib/compat/common/types";
|
||||
import { EVENT_SETTING_SAVED } from "@vrtmrz/livesync-commonlib/compat/events/coreEvents";
|
||||
import { describe, expect, it, vi } from "vitest";
|
||||
import { EVENT_APPLICATION_READY, eventHub } from "@/common/events";
|
||||
import { useReplicationFeature } from "./index";
|
||||
|
||||
type ParseHandler = (documents: PouchDB.Core.ExistingDocument<EntryDoc>[]) => Promise<boolean>;
|
||||
|
||||
function receivedNote(id: string): PouchDB.Core.ExistingDocument<EntryDoc> {
|
||||
return {
|
||||
_id: id,
|
||||
_rev: "1-received",
|
||||
path: `${id}.md`,
|
||||
ctime: 1,
|
||||
mtime: 2,
|
||||
size: 1,
|
||||
children: [],
|
||||
datatype: "plain",
|
||||
type: "plain",
|
||||
eden: {},
|
||||
} as unknown as PouchDB.Core.ExistingDocument<EntryDoc>;
|
||||
}
|
||||
|
||||
function setup() {
|
||||
let applicationReady = false;
|
||||
const settings = {
|
||||
handleFilenameCaseSensitive: false,
|
||||
ignoreFiles: "",
|
||||
maxMTimeForReflectEvents: 0,
|
||||
suspendParseReplicationResult: false,
|
||||
syncIgnoreRegEx: "",
|
||||
syncInternalFiles: false,
|
||||
syncMaxSizeInMB: 0,
|
||||
syncOnlyRegEx: "",
|
||||
useIgnoreFiles: false,
|
||||
};
|
||||
const processSynchroniseResult = vi.fn(async () => true);
|
||||
const settingLoadedHandlers: (() => Promise<boolean>)[] = [];
|
||||
const context = createServiceContext();
|
||||
const keyValueDB = {
|
||||
get: vi.fn(async () => undefined),
|
||||
set: vi.fn(async () => undefined),
|
||||
};
|
||||
const localDatabase = {
|
||||
getRaw: vi.fn(async (id: string) => {
|
||||
if (id === VERSIONING_DOCID) throw { status: 404 };
|
||||
return { _id: id, _rev: "1-received" };
|
||||
}),
|
||||
getDBEntryFromMeta: vi.fn(async (entry: object) => ({ ...entry, data: "received content" })),
|
||||
};
|
||||
let parseHandler: ParseHandler | undefined;
|
||||
const services = {
|
||||
API: { isMobile: vi.fn(() => false), isOnline: true },
|
||||
appLifecycle: {
|
||||
getUnresolvedMessages: { addHandler: vi.fn() },
|
||||
isReady: () => applicationReady,
|
||||
isSuspended: vi.fn(() => false),
|
||||
onSettingLoaded: {
|
||||
addHandler: vi.fn((handler: () => Promise<boolean>) => settingLoadedHandlers.push(handler)),
|
||||
},
|
||||
},
|
||||
context,
|
||||
database: { isDatabaseReady: vi.fn(() => true) },
|
||||
databaseEvents: { onDatabaseInitialised: { addHandler: vi.fn() } },
|
||||
keyValueDB: { kvDB: keyValueDB },
|
||||
localDatabase,
|
||||
path: { getPath: vi.fn((entry: { path: string }) => entry.path) },
|
||||
replication: {
|
||||
databaseQueueCount: { value: 0 },
|
||||
storageApplyingCount: { value: 0 },
|
||||
replicationResultCount: { value: 0 },
|
||||
onBeforeReplicate: { addHandler: vi.fn() },
|
||||
onPrepareCentralRemoteReplication: { addHandler: vi.fn() },
|
||||
onReplicationFailed: { addHandler: vi.fn() },
|
||||
parseSynchroniseResult: {
|
||||
addHandler: vi.fn((handler: ParseHandler) => {
|
||||
parseHandler = handler;
|
||||
}),
|
||||
},
|
||||
processOptionalSynchroniseResult: vi.fn(async () => false),
|
||||
processSynchroniseResult,
|
||||
processVirtualDocument: vi.fn(async () => false),
|
||||
replicateUnattendedByEvent: vi.fn(async () => ({ status: "completed" as const })),
|
||||
},
|
||||
replicator: {
|
||||
createRemoteResource: vi.fn(async () => ({
|
||||
read: vi.fn(async () => new Uint8Array([1])),
|
||||
dispose: vi.fn(),
|
||||
})),
|
||||
onBeforeReplicatorPublication: { addHandler: vi.fn() },
|
||||
onCloseActiveReplication: vi.fn(async () => true),
|
||||
},
|
||||
setting: { currentSettings: vi.fn(() => settings) },
|
||||
tweakValue: {},
|
||||
vault: {
|
||||
isFileSizeTooLarge: vi.fn(() => false),
|
||||
isTargetFile: vi.fn(async () => true),
|
||||
isValidPath: vi.fn(() => true),
|
||||
},
|
||||
};
|
||||
const core = {
|
||||
confirm: {},
|
||||
get localDatabase() {
|
||||
return localDatabase;
|
||||
},
|
||||
rebuilder: {},
|
||||
services,
|
||||
};
|
||||
|
||||
useReplicationFeature(core as never);
|
||||
|
||||
return {
|
||||
context,
|
||||
get applicationReady() {
|
||||
return applicationReady;
|
||||
},
|
||||
processSynchroniseResult,
|
||||
get parseHandler() {
|
||||
return parseHandler;
|
||||
},
|
||||
settingLoadedHandlers,
|
||||
setApplicationReady(value: boolean) {
|
||||
applicationReady = value;
|
||||
},
|
||||
settings,
|
||||
};
|
||||
}
|
||||
|
||||
describe("received change readiness composition", () => {
|
||||
it("applies a queued received document once readiness is established and preserves explicit suspension", async () => {
|
||||
eventHub.offAll();
|
||||
const harness = setup();
|
||||
try {
|
||||
await harness.settingLoadedHandlers[0]();
|
||||
|
||||
await harness.parseHandler!([receivedNote("ready-note")]);
|
||||
await new Promise((resolve) => setTimeout(resolve, 0));
|
||||
expect(harness.processSynchroniseResult).not.toHaveBeenCalled();
|
||||
|
||||
harness.setApplicationReady(true);
|
||||
harness.context.events.emitEvent(EVENT_APPLICATION_READY);
|
||||
harness.context.events.emitEvent(EVENT_APPLICATION_READY);
|
||||
|
||||
await vi.waitFor(() => expect(harness.processSynchroniseResult).toHaveBeenCalledTimes(1));
|
||||
|
||||
harness.settings.suspendParseReplicationResult = true;
|
||||
eventHub.emitEvent(EVENT_SETTING_SAVED, harness.settings as never);
|
||||
harness.setApplicationReady(false);
|
||||
await harness.parseHandler!([receivedNote("suspended-note")]);
|
||||
await new Promise((resolve) => setTimeout(resolve, 0));
|
||||
expect(harness.processSynchroniseResult).toHaveBeenCalledTimes(1);
|
||||
|
||||
harness.setApplicationReady(true);
|
||||
harness.context.events.emitEvent(EVENT_APPLICATION_READY);
|
||||
harness.context.events.emitEvent(EVENT_APPLICATION_READY);
|
||||
await new Promise((resolve) => setTimeout(resolve, 0));
|
||||
expect(harness.processSynchroniseResult).toHaveBeenCalledTimes(1);
|
||||
|
||||
harness.settings.suspendParseReplicationResult = false;
|
||||
eventHub.emitEvent(EVENT_SETTING_SAVED, harness.settings as never);
|
||||
|
||||
await vi.waitFor(() => expect(harness.processSynchroniseResult).toHaveBeenCalledTimes(2));
|
||||
expect(harness.processSynchroniseResult).toHaveBeenLastCalledWith(
|
||||
expect.objectContaining({ _id: "suspended-note", path: "suspended-note.md" })
|
||||
);
|
||||
} finally {
|
||||
eventHub.offAll();
|
||||
}
|
||||
});
|
||||
});
|
||||
@@ -3,7 +3,9 @@ import { createServiceContext } from "@vrtmrz/livesync-commonlib/context";
|
||||
import { VERSIONING_DOCID, type EntryDoc } from "@vrtmrz/livesync-commonlib/compat/common/types";
|
||||
import { REMOTE_FEATURE_GENERATION } from "@vrtmrz/livesync-commonlib/replication";
|
||||
import { promiseWithResolvers } from "octagonal-wheels/promises";
|
||||
import { EVENT_APPLICATION_READY } from "@/common/events";
|
||||
import { useReplicationFeature } from "./index";
|
||||
import { ReplicateResultProcessor } from "./ReplicateResultProcessor";
|
||||
|
||||
type BooleanHandler = (showMessage: boolean) => Promise<boolean>;
|
||||
type ParseHandler = (documents: PouchDB.Core.ExistingDocument<EntryDoc>[]) => Promise<boolean>;
|
||||
@@ -95,6 +97,7 @@ function setup(options: SetupOptions = {}) {
|
||||
return {
|
||||
beforeReplicateHandlers,
|
||||
centralRemoteHandlers,
|
||||
context: services.context,
|
||||
createRemoteResource,
|
||||
dispose,
|
||||
get parseHandler() {
|
||||
@@ -168,6 +171,22 @@ describe("replication serviceFeature composition", () => {
|
||||
expect(createRemoteResource).toHaveBeenCalledOnce();
|
||||
});
|
||||
|
||||
it("continues held results when Commonlib establishes application readiness", () => {
|
||||
const continueHeldDocuments = vi
|
||||
.spyOn(ReplicateResultProcessor.prototype, "continueHeldDocuments")
|
||||
.mockImplementation(() => undefined);
|
||||
try {
|
||||
const { context } = setup();
|
||||
expect(continueHeldDocuments).not.toHaveBeenCalled();
|
||||
|
||||
context.events.emitEvent(EVENT_APPLICATION_READY);
|
||||
|
||||
expect(continueHeldDocuments).toHaveBeenCalledOnce();
|
||||
} finally {
|
||||
continueHeldDocuments.mockRestore();
|
||||
}
|
||||
});
|
||||
|
||||
it("requests owner retirement without awaiting the transition from result application", async () => {
|
||||
const retirement = promiseWithResolvers<boolean>();
|
||||
const onCloseActiveReplication = vi.fn(() => retirement.promise);
|
||||
|
||||
Reference in New Issue
Block a user