Tighten provider adapter boundaries

This commit is contained in:
vorotamoroz
2026-08-30 12:42:34 +00:00
parent 72f033fca4
commit 69cd6250d2
11 changed files with 160 additions and 60 deletions
+1 -7
View File
@@ -15,11 +15,7 @@ import type { StorageAccess } from "@vrtmrz/livesync-commonlib/compat/interfaces
import type { LiveSyncLocalDBEnv } from "@vrtmrz/livesync-commonlib/compat/pouchdb/LiveSyncLocalDB";
import type { LiveSyncCouchDBReplicatorEnv } from "@vrtmrz/livesync-commonlib/compat/replication/couchdb/LiveSyncReplicator";
import type { CheckPointInfo } from "@vrtmrz/livesync-commonlib/compat/replication/journal/JournalSyncTypes";
import type { LiveSyncJournalReplicatorEnv } from "@vrtmrz/livesync-commonlib/compat/replication/journal/LiveSyncJournalReplicatorEnv";
import type {
LiveSyncAbstractReplicator,
LiveSyncReplicatorEnv,
} from "@vrtmrz/livesync-commonlib/compat/replication/LiveSyncAbstractReplicator";
import type { LiveSyncAbstractReplicator } from "@vrtmrz/livesync-commonlib/compat/replication/LiveSyncAbstractReplicator";
import type { ReplicatorInstance } from "@vrtmrz/livesync-commonlib/replication";
import { useTargetFilters } from "@vrtmrz/livesync-commonlib/compat/serviceFeatures/targetFilter";
import { useRemoteConfigurationMigration } from "@vrtmrz/livesync-commonlib/compat/serviceFeatures/remoteConfig";
@@ -51,8 +47,6 @@ export class LiveSyncBaseCore<
>
implements
LiveSyncLocalDBEnv,
LiveSyncReplicatorEnv,
LiveSyncJournalReplicatorEnv,
LiveSyncCouchDBReplicatorEnv,
HasSettings<ObsidianLiveSyncSettings>
{
+51 -8
View File
@@ -19,6 +19,7 @@ import {
type CentralRemoteAdministrationReplicator,
type CentralRemoteAdministrationResult,
type CentralRemoteAdministrationRunner,
type ReplicatorInstance,
type SupportedCapability,
} from "@vrtmrz/livesync-commonlib/replication";
@@ -40,6 +41,19 @@ type CouchDBAdministrationReplicator = CentralRemoteAdministrationReplicator &
type JournalAdministrationClient = Pick<LiveSyncJournalReplicator["client"], "downloadJson">;
function isCentralRemoteAdministrationReplicator(
replicator: ReplicatorInstance
): replicator is CentralRemoteAdministrationReplicator {
return (
"nodeid" in replicator &&
typeof replicator.nodeid === "string" &&
"markRemoteResolved" in replicator &&
typeof replicator.markRemoteResolved === "function" &&
"markRemoteLocked" in replicator &&
typeof replicator.markRemoteLocked === "function"
);
}
async function ensureLocalNodeIdentity(
replicator: CentralRemoteAdministrationReplicator
): Promise<CentralRemoteAdministrationResult | undefined> {
@@ -116,8 +130,12 @@ async function runCentralRemoteAdministration(
function requireCouchDBAdministrationOperations(
replicator: CentralRemoteAdministrationReplicator
): asserts replicator is CouchDBAdministrationReplicator {
const candidate = replicator as Partial<CouchDBAdministrationReplicator>;
if (typeof candidate.connectRemoteCouchDBWithSetting !== "function" || typeof candidate.isMobile !== "function") {
if (
!("connectRemoteCouchDBWithSetting" in replicator) ||
typeof replicator.connectRemoteCouchDBWithSetting !== "function" ||
!("isMobile" in replicator) ||
typeof replicator.isMobile !== "function"
) {
throw new Error("The configured CouchDB administration adapter does not provide milestone access.");
}
}
@@ -164,14 +182,22 @@ function prepareCouchDBMilestoneReader(
};
}
function isJournalAdministrationClient(client: unknown): client is JournalAdministrationClient {
return (
typeof client === "object" &&
client !== null &&
"downloadJson" in client &&
typeof client.downloadJson === "function"
);
}
function requireJournalAdministrationClient(
replicator: CentralRemoteAdministrationReplicator
): JournalAdministrationClient {
const client = (replicator as { readonly client?: JournalAdministrationClient }).client;
if (typeof client?.downloadJson !== "function") {
if (!("client" in replicator) || !isJournalAdministrationClient(replicator.client)) {
throw new Error("The configured Object Storage administration adapter does not provide milestone access.");
}
return client;
return replicator.client;
}
function prepareObjectStorageMilestoneReader(
@@ -191,14 +217,31 @@ function prepareObjectStorageMilestoneReader(
};
}
const runCouchDBCentralRemoteAdministration: CentralRemoteAdministrationRunner = async (replicator, setting, request) =>
await runCentralRemoteAdministration(replicator, setting, request, prepareCouchDBMilestoneReader);
const runCouchDBCentralRemoteAdministration: CentralRemoteAdministrationRunner = async (
replicator,
setting,
request
) => {
if (!isCentralRemoteAdministrationReplicator(replicator)) {
return centralRemoteAdministrationVerificationFailed(
CENTRAL_REMOTE_ADMINISTRATION_FAILURE_REASONS.CAPABILITY_NOT_APPLICABLE
);
}
return await runCentralRemoteAdministration(replicator, setting, request, prepareCouchDBMilestoneReader);
};
const runObjectStorageCentralRemoteAdministration: CentralRemoteAdministrationRunner = async (
replicator,
setting,
request
) => await runCentralRemoteAdministration(replicator, setting, request, prepareObjectStorageMilestoneReader);
) => {
if (!isCentralRemoteAdministrationReplicator(replicator)) {
return centralRemoteAdministrationVerificationFailed(
CENTRAL_REMOTE_ADMINISTRATION_FAILURE_REASONS.CAPABILITY_NOT_APPLICABLE
);
}
return await runCentralRemoteAdministration(replicator, setting, request, prepareObjectStorageMilestoneReader);
};
/** CouchDB mutation and milestone postcondition verification capability. */
export const COUCHDB_CENTRAL_REMOTE_ADMINISTRATION_CAPABILITY: SupportedCapability<CentralRemoteAdministrationRunner> =
+9 -10
View File
@@ -21,7 +21,6 @@ import {
type LiveSyncCouchDBReplicatorEnv,
} from "@vrtmrz/livesync-commonlib/compat/replication/couchdb/LiveSyncReplicator";
import { LiveSyncJournalReplicator } from "@vrtmrz/livesync-commonlib/compat/replication/journal/LiveSyncJournalReplicator";
import type { LiveSyncJournalReplicatorEnv } from "@vrtmrz/livesync-commonlib/compat/replication/journal/LiveSyncJournalReplicatorEnv";
import {
getCouchDBReplicatorConfigurationIdentity,
getObjectStorageReplicatorConfigurationIdentity,
@@ -40,18 +39,19 @@ import {
OBJECT_STORAGE_CENTRAL_REMOTE_ADMINISTRATION_CAPABILITY,
} from "./centralRemoteAdministration";
export type CentralReplicatorProviderHost = LiveSyncCouchDBReplicatorEnv & LiveSyncJournalReplicatorEnv;
/** Host environment sufficient to construct every current central provider. */
export type CentralReplicatorProviderHost = LiveSyncCouchDBReplicatorEnv;
/** Minimal operation required by both central one-shot adapters. */
interface OneShotOutcomeReplicator extends ReplicatorInstance {
openOneShotReplicationWithOutcome(setting: RemoteDBSettings, showResult: boolean): Promise<ReplicationOutcome>;
}
function asOneShotOutcomeReplicator(instance: ReplicatorInstance): OneShotOutcomeReplicator | undefined {
const candidate = instance as Partial<OneShotOutcomeReplicator>;
return typeof candidate.openOneShotReplicationWithOutcome === "function"
? (instance as OneShotOutcomeReplicator)
: undefined;
function isOneShotOutcomeReplicator(instance: ReplicatorInstance): instance is OneShotOutcomeReplicator {
return (
"openOneShotReplicationWithOutcome" in instance &&
typeof instance.openOneShotReplicationWithOutcome === "function"
);
}
async function runOneShotWithOutcome(
@@ -59,11 +59,10 @@ async function runOneShotWithOutcome(
setting: RemoteDBSettings,
showResult: boolean
): Promise<ReplicationOutcome> {
const replicator = asOneShotOutcomeReplicator(instance);
if (!replicator) {
if (!isOneShotOutcomeReplicator(instance)) {
return replicationFailed(new Error("The configured provider does not implement one-shot replication."));
}
return await replicator.openOneShotReplicationWithOutcome(setting, showResult);
return await instance.openOneShotReplicationWithOutcome(setting, showResult);
}
const couchDBUserInitiatedOneShot: UserInitiatedOneShotRunner = async (instance, setting, request) => {
+14 -4
View File
@@ -9,10 +9,10 @@ import {
type LiveSyncCouchDBReplicatorEnv,
} from "@vrtmrz/livesync-commonlib/compat/replication/couchdb/LiveSyncReplicator";
import { LiveSyncJournalReplicator } from "@vrtmrz/livesync-commonlib/compat/replication/journal/LiveSyncJournalReplicator";
import type { LiveSyncJournalReplicatorEnv } from "@vrtmrz/livesync-commonlib/compat/replication/journal/LiveSyncJournalReplicatorEnv";
import { createReplicatorDisposer, snapshotRemoteSettings } from "./shared";
export type ConnectionResourceHost = LiveSyncCouchDBReplicatorEnv & LiveSyncJournalReplicatorEnv;
/** Host environment sufficient to construct either central connection probe. */
export type ConnectionResourceHost = LiveSyncCouchDBReplicatorEnv;
function createCouchDBConnectionProbe(
replicator: LiveSyncCouchDBReplicator,
@@ -60,7 +60,12 @@ function createObjectStorageConnectionProbe(
};
}
/** Build an unpublished CouchDB connection resource for one host. */
/**
* Build an unpublished CouchDB connection probe for one host.
*
* The probe owns both its concrete Replicator and each connection it opens. It
* never publishes that Replicator as the active provider instance.
*/
export function createCouchDBConnectionProbeFactory(host: ConnectionResourceHost): ConnectionProbeFactory {
return (setting) => {
const snapshot = snapshotRemoteSettings(setting);
@@ -68,7 +73,12 @@ export function createCouchDBConnectionProbeFactory(host: ConnectionResourceHost
};
}
/** Build an unpublished Object Storage connection resource for one host. */
/**
* Build an unpublished Object Storage connection probe for one host.
*
* The probe owns its concrete Replicator and never publishes or replaces the
* active provider instance.
*/
export function createObjectStorageConnectionProbeFactory(host: ConnectionResourceHost): ConnectionProbeFactory {
return (setting) => {
const snapshot = snapshotRemoteSettings(setting);
@@ -5,10 +5,10 @@ import {
type LiveSyncCouchDBReplicatorEnv,
} from "@vrtmrz/livesync-commonlib/compat/replication/couchdb/LiveSyncReplicator";
import { LiveSyncJournalReplicator } from "@vrtmrz/livesync-commonlib/compat/replication/journal/LiveSyncJournalReplicator";
import type { LiveSyncJournalReplicatorEnv } from "@vrtmrz/livesync-commonlib/compat/replication/journal/LiveSyncJournalReplicatorEnv";
import { createReplicatorDisposer, snapshotRemoteSettings, type ResourceReplicator } from "./shared";
export type PreferredTweakResourceHost = LiveSyncCouchDBReplicatorEnv & LiveSyncJournalReplicatorEnv;
/** Host environment sufficient to construct either preferred-tweak probe. */
export type PreferredTweakResourceHost = LiveSyncCouchDBReplicatorEnv;
interface PreferredTweakReplicator extends ResourceReplicator {
getRemotePreferredTweakValues(setting: RemoteDBSettings): Promise<RemotePreferredTweakResult>;
@@ -24,7 +24,7 @@ function createPreferredTweakProbe(
};
}
/** Build an unpublished CouchDB preferred-tweak resource for one host. */
/** Build an unpublished, independently disposed CouchDB preferred-tweak probe. */
export function createCouchDBPreferredTweakProbeFactory(host: PreferredTweakResourceHost): PreferredTweakProbeFactory {
return (setting) => {
const snapshot = snapshotRemoteSettings(setting);
@@ -32,7 +32,7 @@ export function createCouchDBPreferredTweakProbeFactory(host: PreferredTweakReso
};
}
/** Build an unpublished Object Storage preferred-tweak resource for one host. */
/** Build an unpublished, independently disposed Object Storage preferred-tweak probe. */
export function createObjectStoragePreferredTweakProbeFactory(
host: PreferredTweakResourceHost
): PreferredTweakProbeFactory {
@@ -5,10 +5,10 @@ import {
type LiveSyncCouchDBReplicatorEnv,
} from "@vrtmrz/livesync-commonlib/compat/replication/couchdb/LiveSyncReplicator";
import { LiveSyncJournalReplicator } from "@vrtmrz/livesync-commonlib/compat/replication/journal/LiveSyncJournalReplicator";
import type { LiveSyncJournalReplicatorEnv } from "@vrtmrz/livesync-commonlib/compat/replication/journal/LiveSyncJournalReplicatorEnv";
import { createReplicatorDisposer, snapshotRemoteSettings, type ResourceReplicator } from "./shared";
export type SecuritySeedResourceHost = LiveSyncCouchDBReplicatorEnv & LiveSyncJournalReplicatorEnv;
/** Host environment sufficient to construct either Security Seed resource. */
export type SecuritySeedResourceHost = LiveSyncCouchDBReplicatorEnv;
/** Minimal private Replicator surface required by a Security Seed resource. */
interface SecuritySeedReplicator extends ResourceReplicator {
@@ -28,12 +28,12 @@ function createSecuritySeedResourceFactory(
};
}
/** Build an unpublished CouchDB Security Seed resource for one host. */
/** Build an unpublished, independently disposed CouchDB Security Seed resource. */
export function createCouchDBSecuritySeedResourceFactory(host: SecuritySeedResourceHost): SecuritySeedResourceFactory {
return createSecuritySeedResourceFactory(() => new LiveSyncCouchDBReplicator(host));
}
/** Build an unpublished Object Storage Security Seed resource for one host. */
/** Build an unpublished, independently disposed Object Storage Security Seed resource. */
export function createObjectStorageSecuritySeedResourceFactory(
host: SecuritySeedResourceHost
): SecuritySeedResourceFactory {
+7
View File
@@ -1,5 +1,12 @@
import type { RemoteDBSettings } from "@vrtmrz/livesync-commonlib/compat/common/types";
/**
* Closeable surface of a concrete Replicator owned by one private resource.
*
* It deliberately exposes no active-provider controls: the resource may use
* the helper for one bounded operation, then must dispose it without
* publishing or replacing the active Replicator.
*/
export interface ResourceReplicator {
closeReplication(): void | Promise<void>;
}
@@ -6,7 +6,12 @@ import {
import { checkSyncInfo } from "@vrtmrz/livesync-commonlib/compat/pouchdb/negotiation";
import { createReplicatorDisposer, snapshotRemoteSettings } from "./shared";
/** Build an owned CouchDB synchronisation-information verifier for one host. */
/**
* Build an unpublished CouchDB synchronisation-information verifier.
*
* The resource owns its concrete Replicator and connection, and cannot replace
* the active provider instance.
*/
export function createCouchDBSynchronisationInformationResourceFactory(
host: LiveSyncCouchDBReplicatorEnv
): SynchronisationInformationResourceFactory {
@@ -17,6 +17,7 @@ import { serialized } from "octagonal-wheels/concurrency/lock_v2";
import { arrayToChunkedArray } from "octagonal-wheels/collection";
import { EVENT_ANALYSE_DB_USAGE, EVENT_REQUEST_PERFORM_GC_V3, eventHub } from "@/common/events";
import type { LiveSyncCouchDBReplicator } from "@vrtmrz/livesync-commonlib/compat/replication/couchdb/LiveSyncReplicator";
import type { ReplicatorInstance } from "@vrtmrz/livesync-commonlib/replication";
import { delay } from "@vrtmrz/livesync-commonlib/compat/common/utils";
import { isNotFoundError } from "@vrtmrz/livesync-commonlib/compat/common/utils.doc";
import { ensureLocalDatabaseMaintenancePrerequisites } from "./maintenancePrerequisites";
@@ -29,6 +30,31 @@ type NoteDocumentID = DocumentID;
type Rev = string;
type ChunkUsageMap = Map<NoteDocumentID, Map<Rev, Set<ChunkID>>>;
type CouchDBCompactionReplicator = ReplicatorInstance &
Pick<LiveSyncCouchDBReplicator, "connectRemoteCouchDBWithSetting">;
type CouchDBGarbageCollectionReplicator = ReplicatorInstance &
Pick<LiveSyncCouchDBReplicator, "getConnectedDeviceList" | "openOneShotReplication">;
function canCompactCouchDBRemote(replicator: ReplicatorInstance): replicator is CouchDBCompactionReplicator {
return (
"connectRemoteCouchDBWithSetting" in replicator &&
typeof replicator.connectRemoteCouchDBWithSetting === "function"
);
}
function canRunCouchDBGarbageCollection(
replicator: ReplicatorInstance
): replicator is CouchDBGarbageCollectionReplicator {
return (
"getConnectedDeviceList" in replicator &&
typeof replicator.getConnectedDeviceList === "function" &&
"openOneShotReplication" in replicator &&
typeof replicator.openOneShotReplication === "function"
);
}
export class LocalDatabaseMaintenance extends LiveSyncCommands {
onunload(): void {
// NO OP.
@@ -737,8 +763,8 @@ Success: ${successCount}, Errored: ${errored}`;
}
async compactDatabase() {
const replicator = this.core.replicator as Partial<LiveSyncCouchDBReplicator>;
if (typeof replicator?.connectRemoteCouchDBWithSetting !== "function") return;
const replicator = this.core.replicator;
if (!canCompactCouchDBRemote(replicator)) return;
const remote = await replicator.connectRemoteCouchDBWithSetting(this.settings, false, false, true);
if (!remote) {
this._notice("Failed to connect to remote for compaction.", "gc-compact");
@@ -841,13 +867,8 @@ Success: ${successCount}, Errored: ${errored}`;
// }
// }
async gcv3() {
const replicator = this.core.replicator as Partial<LiveSyncCouchDBReplicator>;
if (
this.settings.remoteType !== REMOTE_COUCHDB ||
typeof replicator?.openOneShotReplication !== "function" ||
typeof replicator.getConnectedDeviceList !== "function"
)
return;
const replicator = this.core.replicator;
if (this.settings.remoteType !== REMOTE_COUCHDB || !canRunCouchDBGarbageCollection(replicator)) return;
if (!(await this.ensureAvailable("Garbage Collection"))) return;
// Start one-shot replication to ensure all changes are synced before GC.
const r0 = await replicator.openOneShotReplication(this.settings, false, false, "sync");
@@ -5,6 +5,7 @@ import { balanceChunkPurgedDBs, purgeUnreferencedChunks } from "@vrtmrz/livesync
import { LiveSyncCouchDBReplicator } from "@vrtmrz/livesync-commonlib/compat/replication/couchdb/LiveSyncReplicator";
import {
CENTRAL_COMPATIBILITY_REJECTION_REASONS,
type ReplicatorInstance,
type ReplicationFailureRequest,
} from "@vrtmrz/livesync-commonlib/replication";
import { $msg } from "@/common/translation";
@@ -23,6 +24,25 @@ interface CentralCompatibilityRecoveryContext {
readonly services: CentralCompatibilityRecoveryServices;
}
interface PreferredRemoteTweakWriter extends ReplicatorInstance {
setPreferredRemoteTweakSettings(setting: ObsidianLiveSyncSettings): Promise<void>;
}
interface ResolvedRemoteWriter extends ReplicatorInstance {
markRemoteResolved(setting: ObsidianLiveSyncSettings): Promise<void>;
}
function canSetPreferredRemoteTweakSettings(replicator: ReplicatorInstance): replicator is PreferredRemoteTweakWriter {
return (
"setPreferredRemoteTweakSettings" in replicator &&
typeof replicator.setPreferredRemoteTweakSettings === "function"
);
}
function canMarkRemoteResolved(replicator: ReplicatorInstance): replicator is ResolvedRemoteWriter {
return "markRemoteResolved" in replicator && typeof replicator.markRemoteResolved === "function";
}
/**
* Compose central compatibility recovery around the exact failed publication.
* Remote mutations re-admit that publication and become no-ops after a
@@ -126,11 +146,8 @@ Even if you choose to clean up, you will see this option again if you exit Obsid
let updated = false;
await context.services.replicator.runWithActiveReplicatorContext(async (activeContext) => {
if (activeContext !== failedContext) return;
const candidate = activeContext.replicator as typeof activeContext.replicator & {
setPreferredRemoteTweakSettings?: (setting: ObsidianLiveSyncSettings) => Promise<void>;
};
if (typeof candidate.setPreferredRemoteTweakSettings !== "function") return;
await candidate.setPreferredRemoteTweakSettings({ ...effectiveSetting });
if (!canSetPreferredRemoteTweakSettings(activeContext.replicator)) return;
await activeContext.replicator.setPreferredRemoteTweakSettings({ ...effectiveSetting });
updated = true;
});
return updated;
@@ -177,11 +194,8 @@ Even if you choose to clean up, you will see this option again if you exit Obsid
let unlocked = false;
await context.services.replicator.runWithActiveReplicatorContext(async (activeContext) => {
if (activeContext !== failedContext) return;
const replicator = activeContext.replicator as typeof activeContext.replicator & {
markRemoteResolved(setting: ObsidianLiveSyncSettings): Promise<void>;
};
if (typeof replicator.markRemoteResolved !== "function") return;
await replicator.markRemoteResolved(setting);
if (!canMarkRemoteResolved(activeContext.replicator)) return;
await activeContext.replicator.markRemoteResolved(setting);
unlocked = true;
});
if (unlocked) {
+10 -3
View File
@@ -13,6 +13,12 @@ type LocalApplicationActivityOwner = {
runBoundedLocalApplicationActivity<T>(task: () => T | PromiseLike<T>, options?: { label?: string }): Promise<T>;
};
function ownsLocalApplicationActivity(value: object): value is LocalApplicationActivityOwner {
return (
"runBoundedLocalApplicationActivity" in value && typeof value.runBoundedLocalApplicationActivity === "function"
);
}
/**
* Compose result application, automatic triggers, preflight, and central
* compatibility recovery around the existing typed Services.
@@ -28,8 +34,9 @@ export function useReplicationFeature<TContext extends ServiceContext, TCommands
const { services } = core;
// Obsidian adds an application-activity owner to its ReplicatorService.
// Generic hosts retain the former direct-execution fallback.
const localApplicationActivityOwner = services.replicator as typeof services.replicator &
Partial<LocalApplicationActivityOwner>;
const localApplicationActivityOwner = ownsLocalApplicationActivity(services.replicator)
? services.replicator
: undefined;
const resultProcessor = new ReplicateResultProcessor({
currentSettings: () => services.setting.currentSettings(),
keyValueDB: services.keyValueDB.kvDB,
@@ -40,7 +47,7 @@ export function useReplicationFeature<TContext extends ServiceContext, TCommands
fireAndForget(() => services.replicator.onCloseActiveReplication());
},
runLocalApplicationActivity: async (task, options) =>
localApplicationActivityOwner.runBoundedLocalApplicationActivity
localApplicationActivityOwner
? await localApplicationActivityOwner.runBoundedLocalApplicationActivity(task, options)
: await task(),
services: {